"""RAG 하이브리드 검색 — 페르소나 메모리 / 평가 근거 / 정적 지식 KB. 근거: MEMORY_KNOWLEDGE_PERSONA_DESIGN.md §3.6·§4.3 (KB 4-튜플 정책) + MASTERPLAN §3.6 스택(확정 전제, 바꾸지 않음): BGE-M3(dense+sparse, 단일모델) + pgvector(HNSW cosine) + 하이브리드 RRF/가중합 + Anthropic Contextual Retrieval(청크 앞 맥락 프리픽스) + (옵션)BGE-reranker-v2-m3. 설계 불변식: - 정보비대칭은 *프롬프트가 아니라 DB WHERE* 가 강제한다(visible_to[] + sensitivity). 여기 모든 쿼리는 `:role = ANY(visible_to) AND sensitivity <= :sens_max` 를 사전필터로 박는다. (설계서 §4.1: "필터 누락"이 아니라 "코드 경로 부재"가 1차 방어 → 호출부가 role 을 못 바꾸게 함수 시그니처가 정책 4-튜플로 고정.) - "KB는 사전이지 판사가 아니다"(§1, 모순 우선순위) — 검색은 보정용 근거만 돌려준다. 이식성/지연 로딩 원칙 (이 파일의 핵심 계약): - FlagEmbedding/torch/asyncpg-vector 같은 *무거운 의존성*은 **함수 내부에서만 import** 한다. → rag 의존성 미설치 환경에서도 `import app.main` / `import app.services.rag` 가 통과. - 모델/DB 가 없으면 크래시 대신 명확한 `NotConfigured` 를 던진다(라우트가 503 으로 변환). - SQL(pgvector `<=>` 연산자)·인터페이스는 정확히 작성 → DB·모델만 붙으면 그대로 동작. """ from __future__ import annotations import asyncio import json import threading import time from dataclasses import dataclass, field from enum import Enum from typing import TYPE_CHECKING, Any, Optional, Sequence if TYPE_CHECKING: # 타입 체커 전용 — 런타임 import 아님(이식성 유지) import asyncpg # ════════════════════════════════════════════════════════════════════════════ # 0. 예외 — 미구성(모델/DB 부재)은 크래시가 아니라 명시적 신호 # ════════════════════════════════════════════════════════════════════════════ class NotConfigured(RuntimeError): """RAG 구성요소(임베딩 모델 / DB 풀 / 확장)가 준비 안 됨. 라우트가 503(Service Unavailable)로 변환한다. detail 에 무엇이 빠졌는지 명시. """ # ════════════════════════════════════════════════════════════════════════════ # 1. 정책 4-튜플 (설계서 §4.3) — 사전필터 + 회수가중치 + 리랭킹목표 + 주입방식 # 3-AI 가 *같은 물리 테이블 kb.chunk* 를 다른 정책으로 검색한다. # ════════════════════════════════════════════════════════════════════════════ class AIRole(str, Enum): """검색 주체(정보비대칭 축). deps.AIView 와 값 정합(client/counselor/evaluator).""" CLIENT = "client" # 내담자 AI — 본문 비노출(행동단서만), sensitivity<=1 COUNSELOR = "counselor" # 상담사 AI(보조) — DSM diagnostic 차단, sensitivity=0 EVALUATOR = "evaluator" # 평가 AI — taxonomy 정답·논평 포함, sensitivity<=2 @dataclass(frozen=True, slots=True) class RetrievalPolicy: """정책 4-튜플. role 별로 고정(호출부가 임의로 못 푸는 화이트리스트). 설계서 §4.3 표: | 차원 | client | counselor | evaluator | | 사전필터 | diag,theory,tech | theory,tech,micro,ko | (전체) | | | sens<=1 | sens=0 | sens<=2 | | 회수가중치 | dense .7/sparse.3 | dense .5/sparse .5 | sparse .6/dense .4 | | 주입방식 | 본문 비노출(단서) | 본문+예시 | 본문+label_id+bias | """ role: AIRole kinds: tuple[str, ...] # kb_kind 화이트리스트 ('()' = 전체 허용) sens_max: int # sensitivity 상한(이하만 회수) w_dense: float # dense 가중치 w_sparse: float # sparse(BM25/tsvector) 가중치 expose_body: bool # True=본문 주입 / False=행동단서 요약만(내담자 M6) include_label: bool # True=label_id·meta.bias_weight 동봉(평가 채점용) @property def policy_name(self) -> str: """retrieval_log.policy 적재용 정책명(4-튜플 식별).""" return f"{self.role.value}:k={'|'.join(self.kinds) or 'all'}:s<={self.sens_max}" # role → 정책 (설계서 §4.3 SoT). kinds=() 는 "전체 kb_kind 허용"(평가 AI). POLICIES: dict[AIRole, RetrievalPolicy] = { AIRole.CLIENT: RetrievalPolicy( role=AIRole.CLIENT, kinds=("diagnostic", "theory", "technique"), sens_max=1, w_dense=0.7, w_sparse=0.3, expose_body=False, # 본문 비노출 — 행동단서만(R4/M6) include_label=False, ), AIRole.COUNSELOR: RetrievalPolicy( role=AIRole.COUNSELOR, kinds=("theory", "technique", "microskill", "ko_context"), sens_max=0, w_dense=0.5, w_sparse=0.5, expose_body=True, # 본문+예시 include_label=False, ), AIRole.EVALUATOR: RetrievalPolicy( role=AIRole.EVALUATOR, kinds=(), # 전체 kb_kind (taxonomy·supervisor_pattern 정답 포함) sens_max=2, # 평가전용(2)까지. 원천격리(3)는 절대 미회수 w_dense=0.4, w_sparse=0.6, # 라벨명·논평 → sparse 비중↑ expose_body=True, include_label=True, # label_id + meta.bias_weight 동봉(채점 기준) ), } # ════════════════════════════════════════════════════════════════════════════ # 2. 검색 결과 모델 # ════════════════════════════════════════════════════════════════════════════ @dataclass(slots=True) class RetrievedChunk: """회수된 청크 1건. 라우트/오케스트레이터가 system L2 주입에 사용. expose_body=False(내담자) 정책이면 body 는 None, behavior_cue 만 채워 보낸다(M6). """ chunk_id: int score: float # 융합/리랭킹 후 최종 점수 kb_kind: str heading_path: Optional[str] = None context_prefix: Optional[str] = None # Contextual Retrieval 프리픽스 body: Optional[str] = None # chunk_text (expose_body=True 일 때만) behavior_cue: Optional[str] = None # 본문 비노출 정책의 "행동단서 요약" label_id: Optional[int] = None # 평가 정책에서만 meta: dict[str, Any] = field(default_factory=dict) source_id: Optional[str] = None dense_score: float = 0.0 sparse_score: float = 0.0 @dataclass(slots=True) class RetrievalResult: """검색 1회 결과 + 감사 메타(retrieval_log 적재 입력).""" chunks: list[RetrievedChunk] policy_name: str top1_score: float # CRAG 게이트(임계 미달 → 관찰 프레이밍 F-06) latency_ms: int query_text: str degraded: bool = False # reranker/embed fallback 여부(투명성) # CRAG 게이트 임계값(설계서 §3.7 top1_score). 미달이면 호출부가 "관찰 프레이밍"으로 다운그레이드. # 주: 임의 가정값 — Phase 3 파일럿에서 분포 측정 후 확정(M14, "검증됨" 금지). CRAG_TOP1_THRESHOLD = 0.35 # ════════════════════════════════════════════════════════════════════════════ # 3. 임베딩 (BGE-M3) — 지연 로딩 싱글톤 # 무거운 모델/torch import 를 함수 내부로 가둬 import-time 이식성 보장. # ════════════════════════════════════════════════════════════════════════════ _EMBEDDER: Any = None # FlagEmbedding.BGEM3FlagModel 인스턴스(지연 로딩 캐시) _EMBEDDER_FAILED = False # 모델 로드 실패 1회 기록(반복 시도 방지) _EMBEDDER_LOCK = threading.RLock() BGE_M3_MODEL = "BAAI/bge-m3" EMBED_DIM = 1024 # kb.chunk.embedding vector(1024) 와 정합 — 어기면 DB 캐스트 실패 def _get_embedder() -> Any: """BGE-M3 모델 지연 로딩(프로세스 1회). 미설치/로드실패 시 NotConfigured. ⚠️ FlagEmbedding/torch 는 *여기서만* import (모듈 top-level import 금지 — 이식성). """ global _EMBEDDER, _EMBEDDER_FAILED if _EMBEDDER is not None: return _EMBEDDER with _EMBEDDER_LOCK: if _EMBEDDER is not None: return _EMBEDDER if _EMBEDDER_FAILED: raise NotConfigured("BGE-M3 embedder unavailable (prior load failure)") try: from FlagEmbedding import BGEM3FlagModel # 무거운 의존성 — 지연 import except Exception as e: # ImportError 포함(미설치 환경) _EMBEDDER_FAILED = True raise NotConfigured( "FlagEmbedding(BGE-M3) not installed — requirements-rag.txt 필요" ) from e try: # use_fp16: GPU 시 절반정밀(속도). CPU 면 무시됨. _EMBEDDER = BGEM3FlagModel(BGE_M3_MODEL, use_fp16=True) except Exception as e: _EMBEDDER_FAILED = True raise NotConfigured(f"BGE-M3 model load failed: {e}") from e return _EMBEDDER @dataclass(slots=True) class EmbeddedQuery: """질의의 dense+sparse 표현(BGE-M3 단일 호출 산출).""" dense: list[float] # vector(1024) sparse: dict[str, float] = field(default_factory=dict) # {token_id: weight} def embed_query(text: str) -> EmbeddedQuery: """질의 임베딩(dense+sparse). 모델 미가용 시 NotConfigured. 인덱싱 시점(오프라인 배치)에도 같은 함수로 청크 임베딩을 산출한다(동일 모델 재사용). """ with _EMBEDDER_LOCK: model = _get_embedder() out = model.encode( [text], return_dense=True, return_sparse=True, return_colbert_vecs=False, # 멀티벡터는 런타임 회수에 미사용(인덱싱만) ) dense_vec = out["dense_vecs"][0] # numpy → list[float] (asyncpg pgvector 텍스트 캐스트 호환). tolist() 있으면 사용. dense = dense_vec.tolist() if hasattr(dense_vec, "tolist") else list(dense_vec) lexical = out.get("lexical_weights", [{}]) sparse_raw = lexical[0] if lexical else {} sparse = {str(k): float(v) for k, v in dict(sparse_raw).items()} return EmbeddedQuery(dense=dense, sparse=sparse) def _vector_literal(vec: Sequence[float]) -> str: """list[float] → pgvector 텍스트 리터럴 '[a,b,c]'. db.py 가 vector 바이너리 코덱을 아직 안 붙였으므로(주석 TODO), 텍스트 캐스트로 보낸다: `$1::vector` (asyncpg 가 문자열을 vector 로 캐스트). register_vector 코덱이 추후 붙으면 list 직접 바인딩으로 교체 가능. """ if len(vec) != EMBED_DIM: raise NotConfigured( f"embedding dim mismatch: got {len(vec)}, expected {EMBED_DIM} (vector(1024))" ) return "[" + ",".join(f"{x:.7g}" for x in vec) + "]" # ════════════════════════════════════════════════════════════════════════════ # 4. 하이브리드 검색 SQL (pgvector <=> cosine + tsvector BM25, 정책 파라미터화) # 설계서 §4.3 SQL 계약을 03_kb.sql 실제 컬럼명에 정확히 맞춤. # ════════════════════════════════════════════════════════════════════════════ # 파라미터: # $1 = q_dense (vector, '[...]'::vector) # $2 = q_text (tsquery 원문 — websearch_to_tsquery('simple', $2)) # $3 = kinds (text[] — 빈 배열이면 전체 허용) # $4 = role (text — 'client'|'counselor'|'evaluator') # $5 = sens_max(int) # $6 = w_dense (real) # $7 = w_sparse(real) # $8 = pre_k (int — dense/sparse 각 후보 수, 보통 50) # $9 = k (int — 융합 후 반환 수) # # 정보비대칭 강제: 두 CTE 모두 `$4 = ANY(visible_to) AND sensitivity <= $5` 사전필터. # kinds 빈 배열 처리: cardinality($3)=0 이면 kb_kind 조건을 통과(전체). _HYBRID_SQL = """ WITH params AS ( SELECT $1::vector AS q_dense, CASE WHEN $2 = '' THEN NULL ELSE websearch_to_tsquery('simple', $2) END AS q_ts ), dense AS ( SELECT c.chunk_id, 1 - (c.embedding <=> p.q_dense) AS s_dense FROM kb.chunk c, params p WHERE c.embedding IS NOT NULL AND (cardinality($3::text[]) = 0 OR c.kb_kind = ANY($3::text[])) AND $4 = ANY(c.visible_to) AND c.sensitivity <= $5 ORDER BY c.embedding <=> p.q_dense LIMIT $8 ), sparse AS ( SELECT c.chunk_id, ts_rank_cd(to_tsvector('simple', c.chunk_text), p.q_ts) AS s_sparse FROM kb.chunk c, params p WHERE p.q_ts IS NOT NULL AND to_tsvector('simple', c.chunk_text) @@ p.q_ts AND (cardinality($3::text[]) = 0 OR c.kb_kind = ANY($3::text[])) AND $4 = ANY(c.visible_to) AND c.sensitivity <= $5 ORDER BY s_sparse DESC LIMIT $8 ), fused AS ( SELECT COALESCE(d.chunk_id, s.chunk_id) AS chunk_id, COALESCE(d.s_dense, 0) AS s_dense, COALESCE(s.s_sparse, 0) AS s_sparse, COALESCE(d.s_dense, 0) * $6 + COALESCE(s.s_sparse, 0) * $7 AS fused_score FROM dense d FULL OUTER JOIN sparse s USING (chunk_id) ) SELECT c.chunk_id, c.kb_kind, c.heading_path, c.chunk_text, c.context_prefix, c.label_id, c.meta, c.source_id, f.s_dense, f.s_sparse, f.fused_score FROM fused f JOIN kb.chunk c USING (chunk_id) ORDER BY f.fused_score DESC LIMIT $9 """ _PRE_K = 50 # dense/sparse 각 후보 수(설계서 top-50 → reranker top-5) # ════════════════════════════════════════════════════════════════════════════ # 5. Contextual Retrieval 훅 — 청크 앞 맥락 프리픽스 # ════════════════════════════════════════════════════════════════════════════ def apply_contextual_prefix(chunk_text: str, context_prefix: Optional[str]) -> str: """Anthropic Contextual Retrieval: 청크 앞에 1문장 맥락 프리픽스를 결합. 인덱싱 시점에 doc 전체 맥락으로 생성된 context_prefix(kb.chunk.context_prefix)를 *임베딩·BM25 색인 대상*으로는 prefix+body 결합본을 쓰되, **LLM 주입 본문(chunk_text)에는 포함하지 않는다**(03_kb.sql 주석: "표시·LLM 주입(prefix 미포함)"). 이 함수는 인덱싱(색인 텍스트 조립) 경로에서 호출되는 훅. 런타임 회수는 본문만 노출. """ if not context_prefix: return chunk_text return f"{context_prefix.strip()}\n\n{chunk_text}" def _behavior_cue(chunk_text: str, max_len: int = 120) -> str: """본문 비노출 정책(내담자 AI)용 '행동단서 요약'(M6). DSM/이론 본문을 그대로 노출하면 CCD 메타 누설 위험 → 행동 단서만 짧게. 엄밀한 환언은 인덱싱 시점 LLM 이 meta.behavior_cue 로 미리 생성하는 게 이상적이나, 여기선 안전 폴백으로 본문 앞부분만 절단(노출 최소화). meta.behavior_cue 있으면 그걸 우선. """ head = chunk_text.strip().replace("\n", " ") if len(head) <= max_len: return head return head[:max_len].rstrip() + "…" # ════════════════════════════════════════════════════════════════════════════ # 6. 리랭커 (BGE-reranker-v2-m3) — 옵션, 지연 로딩 # ════════════════════════════════════════════════════════════════════════════ _RERANKER: Any = None _RERANKER_FAILED = False RERANKER_MODEL = "BAAI/bge-reranker-v2-m3" def _get_reranker() -> Any: """리랭커 지연 로딩. 미가용이면 None 반환(검색은 융합점수로 폴백 — 크래시 X).""" global _RERANKER, _RERANKER_FAILED if _RERANKER is not None: return _RERANKER if _RERANKER_FAILED: return None try: from FlagEmbedding import FlagReranker # 무거운 의존성 — 지연 import _RERANKER = FlagReranker(RERANKER_MODEL, use_fp16=True) except Exception: _RERANKER_FAILED = True return None return _RERANKER def _rerank( query: str, chunks: list[RetrievedChunk], top_k: int ) -> tuple[list[RetrievedChunk], bool]: """BGE-reranker-v2-m3 로 top_k 재정렬. 미가용 시 (입력 그대로 절단, degraded=True). 리랭킹 점수는 chunk.score 로 덮어쓴다(CRAG top1 게이트가 이 점수를 본다). """ if not chunks: return [], False reranker = _get_reranker() if reranker is None: return chunks[:top_k], True # 폴백: 융합점수 순서 유지 # 본문 비노출 정책이면 behavior_cue, 아니면 body 로 점수화(없으면 prefix). pairs = [ [query, (c.body or c.behavior_cue or c.context_prefix or "")] for c in chunks ] try: scores = reranker.compute_score(pairs, normalize=True) except Exception: return chunks[:top_k], True if not isinstance(scores, (list, tuple)): scores = [scores] for c, s in zip(chunks, scores): c.score = float(s) chunks.sort(key=lambda c: c.score, reverse=True) return chunks[:top_k], False # ════════════════════════════════════════════════════════════════════════════ # 7. retrieval_log 적재 훅 (감사 / 재귀학습, 설계서 §3.7) # ════════════════════════════════════════════════════════════════════════════ _LOG_SQL = """ INSERT INTO kb.retrieval_log (session_id, turn_id, ai_role, query_text, policy, hit_chunk_ids, rerank_scores, top1_score, used_in_answer, latency_ms) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) """ async def log_retrieval( conn: "asyncpg.Connection", *, result: RetrievalResult, ai_role: str, session_id: Optional[str] = None, turn_id: Optional[str] = None, used_in_answer: Optional[bool] = None, ) -> None: """검색 1회를 kb.retrieval_log 에 적재(감사·환각측정·캐시검증). DB/로그 실패는 검색 자체를 막지 않는다(best-effort) — 호출부가 예외를 삼킬 수 있게 여기선 던지되, 라우트가 try/except 로 감싼다. """ hit_ids = [c.chunk_id for c in result.chunks] scores = [c.score for c in result.chunks] await conn.execute( _LOG_SQL, session_id, turn_id, ai_role, result.query_text, result.policy_name, hit_ids, scores, result.top1_score, used_in_answer, result.latency_ms, ) # ════════════════════════════════════════════════════════════════════════════ # 8. 코어 검색 — search_kb (정책 4-튜플로 분기) # ════════════════════════════════════════════════════════════════════════════ async def search_kb( conn: "asyncpg.Connection", *, query: str, role: AIRole, k: int = 5, filters: Optional[dict[str, Any]] = None, rerank: bool = True, pre_k: int = _PRE_K, ) -> RetrievalResult: """정적 지식 KB 하이브리드 검색(dense pgvector cosine + sparse tsvector). Args: conn: asyncpg 커넥션(db.acquire(ai_view=role) 로 RLS 컨텍스트 주입된 것 권장). query: 질의 텍스트(이미 PII 마스킹된 것 — 마스킹은 가드레일 책임). role: AIRole — 정책 4-튜플 선택(사전필터·가중치·노출·라벨). k: 반환 청크 수(리랭킹 후 top-k). filters: 추가 사전필터 — {"kb_kind": [...], "source_id": [...], "sensitivity_max": int}. 정책 화이트리스트를 *좁히는* 방향으로만 적용(넓히지 못함 — 정보비대칭 보존). rerank: BGE-reranker-v2-m3 적용 여부(미설치면 융합점수 폴백, degraded=True). Returns: RetrievalResult (chunks + top1_score + latency + policy_name). Raises: NotConfigured — 임베딩 모델 미가용 OR DB 확장(vector) 미설치. """ t0 = time.perf_counter() policy = POLICIES[role] # (1) 정책 화이트리스트를 filters 로 *좁히기*만 한다(넓히기 금지 — 정보비대칭). kinds = list(policy.kinds) sens_max = policy.sens_max if filters: fk = filters.get("kb_kind") if fk: req = set(fk) if kinds: # 정책이 한정적이면 교집합(좁힘) kinds = [x for x in kinds if x in req] else: # 정책이 전체 허용(평가)이면 요청 그대로(여전히 sensitivity·visible_to 강제) kinds = list(req) fs = filters.get("sensitivity_max") if isinstance(fs, int): sens_max = min(sens_max, fs) # 더 엄격하게만 # (2) 질의 임베딩(dense+sparse). 모델 미가용 → NotConfigured 전파. eq = await asyncio.to_thread(embed_query, query) # CPU 인코딩 → 스레드풀(이벤트루프 비차단) q_dense_lit = _vector_literal(eq.dense) # (3) 하이브리드 SQL 실행. vector 확장 미설치/컬럼 부재면 asyncpg 가 예외 → NotConfigured 변환. try: rows = await conn.fetch( _HYBRID_SQL, q_dense_lit, # $1 ::vector query, # $2 websearch_to_tsquery 원문 kinds, # $3 text[] policy.role.value, # $4 role sens_max, # $5 sens_max policy.w_dense, # $6 policy.w_sparse, # $7 pre_k, # $8 pre_k max(k * 4, k), # $9 융합 후 1차 컷(리랭킹 입력 여유분) ) except Exception as e: # UndefinedFunction(vector 미설치) / UndefinedColumn 등 raise NotConfigured(f"KB hybrid query failed (DB/pgvector not ready): {e}") from e # source_id 추가 좁힘(SQL 후처리 — 화이트리스트 보존, 코드 단순화) src_filter = set(filters.get("source_id", [])) if filters else set() chunks: list[RetrievedChunk] = [] for r in rows: if src_filter and r["source_id"] not in src_filter: continue # asyncpg는 jsonb를 str(JSON text)로 반환 → 파싱. 코덱 등록 시 dict 그대로도 수용. _meta_raw = r["meta"] meta = json.loads(_meta_raw) if isinstance(_meta_raw, str) else dict(_meta_raw or {}) body = r["chunk_text"] if policy.expose_body else None cue = None if not policy.expose_body: # meta.behavior_cue(인덱싱 시 환언) 우선, 없으면 안전 절단 cue = meta.get("behavior_cue") or _behavior_cue(r["chunk_text"]) chunks.append( RetrievedChunk( chunk_id=r["chunk_id"], score=float(r["fused_score"]), kb_kind=r["kb_kind"], heading_path=r["heading_path"], context_prefix=r["context_prefix"], body=body, behavior_cue=cue, label_id=r["label_id"] if policy.include_label else None, meta=meta if policy.include_label else {}, source_id=r["source_id"], dense_score=float(r["s_dense"]), sparse_score=float(r["s_sparse"]), ) ) # (4) 리랭킹(옵션) → top-k degraded = False if rerank and chunks: chunks, degraded = _rerank(query, chunks, k) else: chunks = chunks[:k] top1 = chunks[0].score if chunks else 0.0 latency_ms = int((time.perf_counter() - t0) * 1000) return RetrievalResult( chunks=chunks, policy_name=policy.policy_name, top1_score=top1, latency_ms=latency_ms, query_text=query, degraded=degraded, ) # ════════════════════════════════════════════════════════════════════════════ # 9. 페르소나 메모리 회상 — retrieve_persona_memory (app.turn_embedding, case 스코프) # KB(kb.chunk)와 물리 분리: 에피소드 메모리는 app 스키마. case_id 스코프 강제(M5/cross-trainee). # ════════════════════════════════════════════════════════════════════════════ # 설계서 §2-A 3단계 회상: open_threads/homework 우선쿼리 → 하이브리드 top-50 → reranker top-5. # 내담자 AI 뷰: CCD/정답/평가 *절대 미포함*(app.turn_embedding 에는 표면 발화만). case 스코프가 1차 격리. _MEMORY_SQL = """ WITH params AS ( SELECT $1::vector AS q_dense, CASE WHEN $2 = '' THEN NULL ELSE websearch_to_tsquery('simple', $2) END AS q_ts ), dense AS ( SELECT te.turn_id, te.seq, 1 - (te.dense <=> p.q_dense) AS s_dense FROM app.turn_embedding te, params p WHERE te.case_id = $3::uuid ORDER BY te.dense <=> p.q_dense LIMIT $5 ) SELECT d.turn_id, d.seq, d.s_dense FROM dense d ORDER BY d.s_dense DESC LIMIT $4 """ async def retrieve_persona_memory( conn: "asyncpg.Connection", *, case_id: str, query: str, k: int = 5, pre_k: int = _PRE_K, ) -> RetrievalResult: """회기 시작 episodic recall (내담자 연속성, app.turn_embedding HNSW). case_id 스코프 강제 → 학습자 간 기억 오염 차단(M5/T4). 본문(turn 텍스트)은 호출부(memory.build_recall_context)가 turns 조인으로 가져오되, 여기선 turn_id+점수만 회수(검색 책임 분리). CCD/정답은 이 경로에 *구조적으로* 존재하지 않음(코드경로 부재 1차방어). Raises: NotConfigured — 임베딩 모델/DB 미가용. """ t0 = time.perf_counter() eq = await asyncio.to_thread(embed_query, query) # CPU 인코딩 → 스레드풀(이벤트루프 비차단) q_dense_lit = _vector_literal(eq.dense) try: rows = await conn.fetch( _MEMORY_SQL, q_dense_lit, # $1 ::vector query, # $2 (현재 dense-only; sparse 백필은 turn_embedding.sparse 추가 시 확장) case_id, # $3 ::uuid (case 스코프) k, # $4 pre_k, # $5 ) except Exception as e: raise NotConfigured(f"persona memory query failed (DB/pgvector not ready): {e}") from e chunks = [ RetrievedChunk( chunk_id=0, # turn 은 UUID — chunk_id(int) 의미 없음, meta 에 보관 score=float(r["s_dense"]), kb_kind="episodic", meta={"turn_id": str(r["turn_id"]), "seq": r["seq"]}, dense_score=float(r["s_dense"]), ) for r in rows ] top1 = chunks[0].score if chunks else 0.0 latency_ms = int((time.perf_counter() - t0) * 1000) return RetrievalResult( chunks=chunks, policy_name=f"client:episodic:case={case_id}", top1_score=top1, latency_ms=latency_ms, query_text=query, ) # ════════════════════════════════════════════════════════════════════════════ # 10. 평가 근거 회수 — retrieve_eval_grounding (평가 AI 채점 근거) # evaluator 정책: taxonomy 정답·슈퍼바이저 논평 포함(sensitivity<=2), label_id 동봉. # ════════════════════════════════════════════════════════════════════════════ async def retrieve_eval_grounding( conn: "asyncpg.Connection", *, query: str, k: int = 5, kinds: Optional[Sequence[str]] = None, rerank: bool = True, ) -> RetrievalResult: """평가 AI 채점 근거 회수 — DSM/이론/taxonomy 정답라벨 + 논평. search_kb(role=EVALUATOR) 래퍼. CRAG 게이트(top1_score < 임계 → 관찰 프레이밍, F-06)는 호출부(evaluator)가 result.top1_score 로 판단한다. label_id·meta.bias_weight 동봉. kinds: 평가 차원에 따라 좁히기(예: 기법 채점 → ['technique','supervisor_pattern']). """ filters = {"kb_kind": list(kinds)} if kinds else None return await search_kb( conn, query=query, role=AIRole.EVALUATOR, k=k, filters=filters, rerank=rerank, ) # ════════════════════════════════════════════════════════════════════════════ # 11. 인덱싱 트리거 (관리자 — 오프라인 배치, 런타임 아님) # content_hash 증분(graphify 차용). 실제 청킹/Contextual prefix 생성은 LLM 배치. # ════════════════════════════════════════════════════════════════════════════ @dataclass(slots=True) class IndexRequest: """문서 1건 인덱싱 요청(관리자 엔드포인트 입력).""" source_id: str doc_uri: str chunks: list[dict[str, Any]] # [{seq, heading_path, chunk_text, context_prefix?, visible_to?, sensitivity?, meta?}] version: int = 1 content_hash: Optional[str] = None # 미지정 시 chunk_text 합으로 계산 @dataclass(slots=True) class IndexResult: doc_id: Optional[int] chunks_indexed: int skipped_unchanged: bool # content_hash 동일 → 증분 스킵 embedded: bool # 임베딩 실제 적재 여부(모델 미가용 시 False) degraded: bool = False def _content_hash(chunks: list[dict[str, Any]]) -> str: """청크 본문 합의 SHA256 — 변경감지(증분 인덱싱, kb.document.content_hash).""" import hashlib # 표준 라이브러리 — 지연 불필요하나 일관성 위해 함수 내부 h = hashlib.sha256() for c in chunks: h.update((c.get("chunk_text") or "").encode("utf-8")) return h.hexdigest() async def index_document( conn: "asyncpg.Connection", req: IndexRequest, ) -> IndexResult: """문서 인덱싱(관리자 트리거). content_hash 증분 + 청크 임베딩 적재. 절차(설계서 §3.6 + graphify 증분): 1. content_hash 계산 → kb.document 동일 활성본 있으면 스킵(증분). 2. kb.document UPSERT(신규 version) → doc_id. 3. 각 청크: 임베딩(prefix+body 결합본을 색인 텍스트로 — Contextual Retrieval) → kb.chunk INSERT. 4. 임베딩 모델 미가용 시: embedding NULL 로 적재(텍스트만, BM25 만 동작) + degraded=True. ⚠️ 무거운 작업(임베딩) → 본래는 백그라운드 워커/배치. 라우트는 BackgroundTasks 로 위임 권장. DSM verbatim 저작권(license C/D): source.external_llm_ok=false 가드는 source 등록 시점 책임. Raises: NotConfigured — DB(kb 스키마/vector) 미가용. """ content_hash = req.content_hash or _content_hash(req.chunks) # (1) 증분 — 동일 source/uri/version 활성본의 content_hash 비교 try: existing = await conn.fetchrow( """ SELECT doc_id, content_hash FROM kb.document WHERE source_id = $1 AND doc_uri = $2 AND is_active ORDER BY version DESC LIMIT 1 """, req.source_id, req.doc_uri, ) except Exception as e: raise NotConfigured(f"kb.document not ready (DB/schema): {e}") from e if existing and existing["content_hash"] == content_hash: return IndexResult( doc_id=existing["doc_id"], chunks_indexed=0, skipped_unchanged=True, embedded=False, ) # (2) 새 문서 행 — 기존본 supersede + 신규 active version = req.version if existing: version = max(version, 1) # 호출부가 version 증가 책임(여기선 UNIQUE 충돌 방어만) try: doc = await conn.fetchrow( """ INSERT INTO kb.document (source_id, doc_uri, version, content_hash, is_active, indexed_at) VALUES ($1, $2, $3, $4, TRUE, now()) RETURNING doc_id """, req.source_id, req.doc_uri, version, content_hash, ) except Exception as e: raise NotConfigured(f"kb.document insert failed: {e}") from e doc_id = doc["doc_id"] # 직전 활성본 비활성화(증분 supersede) if existing: await conn.execute( "UPDATE kb.document SET is_active = FALSE, superseded_by = $2 WHERE doc_id = $1", existing["doc_id"], doc_id, ) # (3) 청크 임베딩 + 적재. 모델 미가용 → embedding NULL 폴백(BM25 만). embedded = True degraded = False try: embedder = _get_embedder() except NotConfigured: embedder = None embedded = False degraded = True indexed = 0 for c in req.chunks: chunk_text = c.get("chunk_text") or "" if not chunk_text: continue context_prefix = c.get("context_prefix") emb_lit: Optional[str] = None sparse_json: Optional[str] = None # jsonb 바인딩용 직렬화 문자열(asyncpg는 dict 자동인코딩 안 함) if embedder is not None: # Contextual Retrieval: prefix+body 결합본을 *색인 대상* 으로 임베딩(주입 본문은 body 만). index_text = apply_contextual_prefix(chunk_text, context_prefix) eq = await asyncio.to_thread(embed_query, index_text) # CPU 인코딩 → 스레드풀 emb_lit = _vector_literal(eq.dense) sparse_json = json.dumps(eq.sparse) await conn.execute( """ INSERT INTO kb.chunk (doc_id, source_id, kb_kind, seq, heading_path, chunk_text, context_prefix, embedding, sparse_vec, visible_to, sensitivity, label_id, meta, token_count) VALUES ($1, $2, $3, $4, $5, $6, $7, $8::vector, $9::jsonb, COALESCE($10::text[], ARRAY['client','counselor','evaluator']), COALESCE($11, 0), $12, COALESCE($13::jsonb, '{}'::jsonb), $14) """, doc_id, req.source_id, c.get("kb_kind") or "theory", c.get("seq", indexed), c.get("heading_path"), chunk_text, context_prefix, emb_lit, sparse_json, c.get("visible_to"), c.get("sensitivity"), c.get("label_id"), json.dumps(c.get("meta")) if c.get("meta") is not None else None, c.get("token_count"), ) indexed += 1 return IndexResult( doc_id=doc_id, chunks_indexed=indexed, skipped_unchanged=False, embedded=embedded, degraded=degraded, ) __all__ = [ "NotConfigured", "AIRole", "RetrievalPolicy", "POLICIES", "RetrievedChunk", "RetrievalResult", "CRAG_TOP1_THRESHOLD", "EmbeddedQuery", "embed_query", "apply_contextual_prefix", "search_kb", "retrieve_persona_memory", "retrieve_eval_grounding", "log_retrieval", "IndexRequest", "IndexResult", "index_document", ]