""" admin_service.py ──────────────────────────────────────────── 벡터 DB(Qdrant) 큐레이션 어드민 API + 단일 페이지 웹 UI. 기능: - 검색(semantic): 질문 임베딩 → 벡터 검색 - 목록(browse): 페이지네이션 조회 - 단건 조회 / 삭제 / 일괄 삭제 - 추가: 질문/답변 텍스트 → 임베딩 → 결정적 ID로 upsert 실시간 /ask 서비스(run_service_qa.py)와 같은 모듈/벡터스토어를 공유하되, 별도 프로세스(컨테이너)로 분리해 운영한다. Qdrant 전용 기능을 사용한다. """ import os from datetime import datetime, timezone from pathlib import Path from typing import Optional, List, Dict, Any from urllib.parse import unquote import numpy as np from fastapi import FastAPI, HTTPException from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import HTMLResponse from pydantic import BaseModel, Field from api_clients import TEIEmbeddingClient from vector_store import get_vector_store, VECTOR_STORE, make_point_id app = FastAPI(title="RAG 벡터DB 큐레이션 어드민") # 별도 origin(다른 포트/도메인)의 프론트엔드가 호출할 수 있도록 CORS 허용. # 내부망 도구이므로 기본 전체 허용. 필요 시 ADMIN_CORS_ORIGINS="http://host:port,..." 로 제한. _cors = os.getenv("ADMIN_CORS_ORIGINS", "*") _origins = ["*"] if _cors.strip() == "*" else [o.strip() for o in _cors.split(",") if o.strip()] app.add_middleware( CORSMiddleware, allow_origins=_origins, allow_methods=["*"], allow_headers=["*"], ) DATA_DIR = Path(os.getenv("DATA_DIR", "/app/data")) def _now_iso() -> str: """등록/수정 시각 기본값 (UTC ISO-8601)""" return datetime.now(timezone.utc).isoformat() print("[Admin] 초기화 시작...") embed_client = TEIEmbeddingClient() vector_store = get_vector_store() vector_store.load(str(DATA_DIR)) print(f"[Admin] 벡터 스토어({VECTOR_STORE}) 로드 완료: {vector_store.count()}개") if VECTOR_STORE != "qdrant": print("[Admin] ⚠️ 경고: 어드민의 조회/삭제/추가 기능은 Qdrant에서만 동작합니다 " f"(현재 VECTOR_STORE={VECTOR_STORE}).") else: vector_store.ensure_admin_payload_indexes() # ─────────────────────────────────────────── # 요청/응답 모델 # ─────────────────────────────────────────── class SearchRequest(BaseModel): query: str = Field(..., min_length=1, max_length=500) top_k: int = Field(20, ge=1, le=100) threshold: Optional[float] = None category: Optional[str] = None # 정확일치 필터 source: Optional[str] = None # 정확일치 필터 class AddRequest(BaseModel): q: str = Field(..., min_length=1) a: str = Field(..., min_length=1) category: Optional[str] = None url: Optional[str] = None source: str = "admin_manual" source_id: Optional[str] = None source_created_at: Optional[str] = None class UpdateRequest(BaseModel): """전달된 필드만 수정. q가 바뀌면 재임베딩한다.""" q: Optional[str] = None a: Optional[str] = None category: Optional[str] = None url: Optional[str] = None source_created_at: Optional[str] = None class KeywordSearchRequest(BaseModel): """DB 전체 대상 문자열 포함(부분일치) 검색""" keyword: str = Field(..., min_length=1, max_length=200) field: str = Field("both", pattern="^(both|q|a)$") # 검색 대상 필드 category: Optional[str] = None source: Optional[str] = None skip: int = Field(0, ge=0) limit: int = Field(50, ge=1, le=200) class DeleteBatchRequest(BaseModel): ids: List[str] = Field(..., min_items=1) def _row(item: Dict[str, Any], score: Optional[float] = None) -> Dict[str, Any]: """검색/조회 결과를 UI 친화 형태로 변환""" meta = item.get("meta", {}) or {} row = { "id": item.get("id"), "q": meta.get("q"), "a": meta.get("a"), "category": meta.get("category"), "source": meta.get("source"), "source_id": meta.get("source_id"), "url": meta.get("url"), "source_created_at": meta.get("source_created_at"), "indexed_at": meta.get("indexed_at"), "updated_at": meta.get("updated_at"), } if score is not None: row["score"] = round(float(score), 4) return row def _require_qdrant(): if VECTOR_STORE != "qdrant": raise HTTPException( status_code=400, detail="이 기능은 Qdrant 벡터스토어에서만 지원됩니다. VECTOR_STORE=qdrant 로 실행하세요.", ) # ─────────────────────────────────────────── # API # ─────────────────────────────────────────── @app.get("/api/stats") def stats(): return {"vector_store": VECTOR_STORE, "count": vector_store.count()} def _normalize_query_param(value: Optional[str]) -> Optional[str]: """쿼리 파라미터 이중 URL 인코딩 복구 (Java RestTemplate + toUriString 조합 대응).""" if value is None: return None v = value.strip() if not v: return v for _ in range(3): decoded = unquote(v) if decoded == v: break v = decoded return v def _filter_dict(category: Optional[str], source: Optional[str]) -> Optional[Dict[str, Any]]: f: Dict[str, Any] = {} if category: cat = (_normalize_query_param(category) or category).strip() if cat == "미분류": f["category"] = "__EMPTY__" else: f["category"] = cat if source: src = (_normalize_query_param(source) or source).strip() if src: f["source"] = src return f or None @app.post("/api/search") def search(req: SearchRequest): vecs = embed_client.embed([req.query], normalize=True, is_query=True) if not vecs: raise HTTPException(status_code=502, detail="임베딩 실패") query_vec = np.array(vecs[0], dtype="float32") results = vector_store.search( query_vec, top_k=req.top_k, threshold=req.threshold, filter_dict=_filter_dict(req.category, req.source), ) return {"count": len(results), "items": [_row(r, r.get("score")) for r in results]} @app.post("/api/keyword-search") def keyword_search(req: KeywordSearchRequest): """문자열 포함(부분일치) 검색 — 의미 검색과 달리 키워드가 실제 포함된 항목만 반환""" _require_qdrant() fields = {"both": ("q", "a"), "q": ("q",), "a": ("a",)}[req.field] items, total, exhausted = vector_store.keyword_search( req.keyword, fields=fields, filter_dict=_filter_dict(req.category, req.source), skip=req.skip, limit=req.limit, ) return { "count": len(items), "total": total, "exhausted": exhausted, "items": [_row(it) for it in items], } @app.get("/api/categories") def categories(): """등록된 category 목록 + 건수 (어드민 분류 필터/트리용)""" _require_qdrant() counts = vector_store.distinct_payload_values("category") items = [{"category": k, "count": v} for k, v in sorted(counts.items())] return {"count": len(items), "items": items} @app.get("/api/points") def list_points( limit: int = 50, offset: Optional[str] = None, page: Optional[int] = None, size: Optional[int] = None, category: Optional[str] = None, source: Optional[str] = None, ): _require_qdrant() fd = _filter_dict(category, source) if page is not None: pg = max(0, page) sz = max(1, min(size or 20, 200)) items, total = vector_store.list_points_page(pg, sz, fd) return { "page": pg, "size": sz, "total": total, "items": [_row(it) for it in items], } items, next_offset = vector_store.list_points( limit=limit, offset=offset, filter_dict=fd ) return { "count": len(items), "items": [_row(it) for it in items], "next_offset": next_offset, } @app.get("/api/points/{point_id}") def get_point(point_id: str): _require_qdrant() item = vector_store.get_by_id(point_id) if not item: raise HTTPException(status_code=404, detail="해당 ID의 항목이 없습니다.") return _row(item) @app.put("/api/points/{point_id}") def update_point(point_id: str, req: UpdateRequest): """ 기존 항목 수정. 전달된 필드만 갱신한다. - q 기반 결정적 ID → q 변경 시 새 ID로 upsert 후 옛 ID 삭제 - 추가/업로드와 동일: 더 최신 source_created_at이 이미 있으면 스킵 """ _require_qdrant() existing = vector_store.get_by_id(point_id) if not existing: raise HTTPException(status_code=404, detail="해당 ID의 항목이 없습니다.") meta: Dict[str, Any] = dict(existing.get("meta") or {}) if req.q is not None: meta["q"] = req.q.strip() if req.a is not None: meta["a"] = req.a if req.category is not None: meta["category"] = req.category if req.url is not None: meta["url"] = req.url if req.source_created_at is not None: meta["source_created_at"] = req.source_created_at.strip() or None if not meta.get("q") or not meta.get("a"): raise HTTPException(status_code=400, detail="q와 a는 비울 수 없습니다.") if not meta.get("source_created_at") or not str(meta.get("source_created_at")).strip(): raise HTTPException(status_code=400, detail="source_created_at은 필수입니다.") meta["source_created_at"] = str(meta["source_created_at"]).strip() # 질문 재임베딩 (q가 안 바뀌었어도 일관성을 위해 항상 재임베딩) vecs = embed_client.embed([meta["q"]], normalize=True, is_query=False) if not vecs: raise HTTPException(status_code=502, detail="임베딩 실패") new_id = make_point_id(meta) vectors = np.array([vecs[0]], dtype="float32") if hasattr(vector_store, "upsert_vectors"): stats = vector_store.upsert_vectors(vectors, [meta], skip_if_older=True) if stats.get("skipped"): return { "skipped": True, "reason": stats.get("last_skip_reason") or "older_source_created_at", "id": stats.get("last_id") or new_id, } else: meta.setdefault("indexed_at", _now_iso()) meta["updated_at"] = _now_iso() vector_store.add_vectors(vectors, [meta], skip_if_older=True) # ID가 바뀐 경우(질문 변경) 옛 항목 제거 if new_id != point_id: vector_store.delete_by_id(point_id) # 레거시(source_id 기반) ID로 남아 있는 동일 q 항목 정리 if hasattr(vector_store, "find_by_exact_q"): legacy = vector_store.find_by_exact_q(meta["q"]) if legacy and str(legacy.get("id")) not in (str(new_id), str(point_id)): vector_store.delete_by_id(str(legacy["id"])) return {"updated": new_id, "moved": new_id != point_id, "skipped": False} @app.delete("/api/points/{point_id}") def delete_point(point_id: str): _require_qdrant() if not vector_store.get_by_id(point_id): raise HTTPException(status_code=404, detail="해당 ID의 항목이 없습니다.") vector_store.delete_by_id(point_id) return {"deleted": point_id} @app.post("/api/points/delete-batch") def delete_batch(req: DeleteBatchRequest): _require_qdrant() vector_store.delete_by_ids(req.ids) return {"deleted": len(req.ids)} @app.post("/api/points") def add_point(req: AddRequest): _require_qdrant() if not req.source_created_at or not req.source_created_at.strip(): raise HTTPException(status_code=400, detail="source_created_at은 필수입니다.") # 문서 임베딩(is_query=False) → ingest와 동일 방식 vecs = embed_client.embed([req.q], normalize=True, is_query=False) if not vecs: raise HTTPException(status_code=502, detail="임베딩 실패") source_id = req.source_id or f"manual_{os.urandom(6).hex()}" now = _now_iso() meta: Dict[str, Any] = { "q": req.q.strip(), "a": req.a, "category": req.category, "source": req.source, "source_id": source_id, "url": req.url, "source_created_at": req.source_created_at.strip(), "indexed_at": now, "updated_at": now, } point_id = make_point_id(meta) vectors = np.array([vecs[0]], dtype="float32") if hasattr(vector_store, "upsert_vectors"): stats = vector_store.upsert_vectors(vectors, [meta], skip_if_older=True) if stats.get("skipped"): return { "skipped": True, "reason": stats.get("last_skip_reason") or "older_source_created_at", "id": stats.get("last_id") or point_id, } action = stats.get("last_action") or "inserted" return { "created": stats.get("last_id") or point_id, "source_id": source_id, "action": action, "skipped": False, } vector_store.add_vectors(vectors, [meta], skip_if_older=True) return {"created": point_id, "source_id": source_id, "skipped": False} # ─────────────────────────────────────────── # 웹 UI (빌드 불필요 단일 페이지) # ─────────────────────────────────────────── @app.get("/", response_class=HTMLResponse) def index(): return HTML_PAGE HTML_PAGE = """ 벡터DB 큐레이션 어드민

🗂️ 벡터DB 큐레이션 어드민

로딩...
의미검색
키워드검색
목록
추가

추가 시 질문이 임베딩되어 검색 대상이 됩니다. 같은 항목을 다시 추가하면 덮어쓰기됩니다.

"""