이전 글에서 GitHub App과 Webhook으로 코드 변경 시 자동 인덱싱이 되도록 만들었습니다.
이번 글은 실제로 검색을 요청했을 때 어떤 흐름으로 동작하는지, 검색 API가 어떤 구조로 설계되어 있는지를 써볼까 합니다.
SSE
SSE란?
: SSE는 Server-Sent Events의 약자
: 서버가 클라이언트에게 데이터를 지속적으로 밀어주는(push) 단방향 통신 방식
HTTP 요청과 SSE 차이
▶︎ 일반적인 HTTP 요청
| 클라이언트 → 요청 → 서버 클라이언트 ← 응답 ← 서버 (한 번에 전부) |
▶︎ SSE
| 클라이언트 → 요청 → 서버 클라이언트 ← 데이터 조각 1 ← 서버 클라이언트 ← 데이터 조각 2 ← 서버 클라이언트 ← 데이터 조각 3 ← 서버 ... 클라이언트 ← 완료 신호 ← 서버 |
: 응답이 완성될 때까지 기다리는 것이 아니라, 데이터가 준비되는 즉시 조각조각 전달
: 연결은 서버가 완료 신호를 보낼 때까지 유지
WebSocket과 차이점
| 구분 | SSE | WebSocket |
| 방향 | 서버 → 클라이언트 (단방향) | 양방향 |
| 프로토콜 | HTTP | 별도 프로토콜 (ws://) |
| 재연결 | 브라우저가 자동 처리 | 직접 구현 필요 |
| 적합한 경우 | 서버 푸시, 스트리밍 응답 | 채팅, 실시간 게임 |
: 코드 검색처럼 "요청 한 번 → 서버가 계속 밀어주는" 구조에는 SSE가 더 적합
검색 API 스펙
POST /api/v1/search
Content-Type: application/json
class SearchRequest(BaseModel):
question: str = Field(min_length=1)
top_k: int = Field(default=5, ge=1, le=20)
# null 또는 빈 배열이면 전체 레포 검색
repositories: list[str] | None = None
| 필드 | 타입 | 필수 여부 | 설명 |
| question | string | O | 지연어 질문 |
| top_k | int (1~20) | 선택 (기본 5) | 참조할 코드 청크 수 |
| repositories | string [] | 선택 | 검색 대상 레포 이름 목록 |
- repositories가 null이거나 빈 배열이면 전체 레포를 대상으로 검색
- 응답 Content-Type은 text/event-stream
SSE 이벤트 구조
스트리밍 응답은 4가지 이벤트 타입으로 구성
metadata
: 검색이 완료된 직후 가장 먼저 전송
: LLM이 답변을 생성하기 전에 참조 코드 목록을 먼저 내림
{
"type": "metadata",
"references": [
{
"file_path": "payment-service/src/retry/RetryHandler.java",
"start_line": 42,
"end_line": 87,
"language": "java",
"score": 0.9134,
"snippet": "public class RetryHandler { ...",
"is_deprecated": false,
"warning": null
}
]
}
- 프론트엔드는 LLM이 답변을 생성하는 동안 참조 파일 목록을 미리 화면에 표시할 수 있고 사용자는 어떤 코드를 근거로 답변이 생성되는지 답변이 완성되기 전에 미리 확인할 수 있음
- is_deprecated가 true인 청크는 warning 필드에 안내 메시지가 함께 담겨 옵니다. 오래되거나 더 이상 사용하지 않는 파일을 참조한 경우 프론트엔드가 경고를 표시
token
{ "type": "token", "text": "결제 실패 시 재시도는", "done": false }
{ "type": "token", "text": " RetryHandler 클래스에서 처리합니다.", "done": true }
- LLM이 생성한 텍스트를 토큰 단위로 전달
done
{ "type": "done", "done": true }
- 스트리밍이 정상 종료됐음을 알림
error
{
"type": "error",
"code": "SEARCH_TIMEOUT",
"message": "요청 시간이 초과되었습니다.",
"done": true
}
- 타임아웃이나 내부 오류 발생 시 전송
전체 이벤트 흐름
| metadata (참조 코드 목록) → token × N (LLM 답변 스트리밍) → done |
FastAPI에서 SSE 구현
FastAPI란?
Python으로 만드는 웹 API 서버 프레임워크
FastAPI에서는 StreamingResponse와 AsyncGenerator를 조합해 SSE를 구현
@router.post("")
async def search(
request: SearchRequest,
db: Session = Depends(get_db),
search_service: SearchApplicationService = Depends(get_search_service),
) -> StreamingResponse:
repository_ids = _resolve_repository_ids(db, request.repositories)
return StreamingResponse(
search_service.stream(request, repository_ids=repository_ids),
media_type="text/event-stream",
)
- StreamingResponse에 AsyncGenerator를 넘기면 제너레이터가 값을 yield할 때마다 클라이언트로 즉시 전송
각 SSE 이벤트는 아래 형식으로 직렬화
def _sse_event(payload: dict[str, object]) -> str:
body = json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
return f"data: {body}\n\n"
실제 스트리밍
async def _stream(self, request, repository_ids):
embedding = await self.embedding_service.embed(request.question)
# Dense와 BM25를 병렬 호출
dense_chunks, bm25_chunks = await asyncio.gather(
self.vector_search_service.search(embedding, candidate_k, repository_ids),
self.keyword_search_service.search(request.question, candidate_k, repository_ids),
)
merged = rrf_merge(dense_chunks, bm25_chunks, top_k=candidate_k)
chunks = self.reranker.rerank(request.question, merged)[:request.top_k]
# 검색 결과를 먼저 전송
yield _sse_event({"type": "metadata", "references": [...]})
# LLM 답변을 토큰 단위로 스트리밍
async for chunk in self.ollama_client.generate_stream(prompt):
yield _sse_event({"type": "token", "text": chunk.text, "done": chunk.done})
레포 필터 설계
: repositories 파라미터로 특정 레포만 검색 범위를 좁힐 수 있음
: Dense 검색(Chroma)과 BM25에서 이 필터가 다른 방식으로 동작
Dense 검색 — Chroma where 절
def _build_where(self, repository_ids: list[int] | None) -> dict | None:
if not repository_ids:
return None
if len(repository_ids) == 1:
return {"repository_id": {"$eq": repository_ids[0]}}
return {"repository_id": {"$in": repository_ids}}
- Chroma는 벡터 검색 시 메타데이터 조건을 where 절로 받아 DB 레벨에서 필터링
- 레포가 하나면 $eq, 여러 개면 $in 조건을 사용
- 필터링이 Chroma 내부에서 처리되므로 불필요한 결과가 메모리로 올라오지 않음
BM25 — 인메모리 post-filter
for idx, score in candidates:
if len(results) >= top_k:
break
meta = self._chunk_metas[idx].get("metadata", {})
if repository_ids and meta.get("repository_id") not in repository_ids:
continue
results.append(RetrievedChunk(...))
- BM25는 Chroma 같은 DB가 아니라 인메모리 인덱스
- 검색 조건을 미리 넘길 방법이 없으므로, 전체 후보를 스코어 순으로 정렬한 뒤 repository_id를 직접 확인해 필터링함
필터 시점이 다른 이유
| 구분 | Dense (Chroma) | BM25 |
| 저장 위치 | 외부 벡터 DB | 인메모리 인덱스 |
| 필터 시점 | 검색 전 (DB 레벨) | 검색 후 (결과 순회) |
| 필터 방식 | where 절 | repository_id 직접 비교 |
두 방식 모두 최종적으로 동일한 레포 범위를 대상으로 검색하고, RRF로 합산
레포 이름 → ID 변환
def _resolve_repository_ids(db: Session, repositories: list[str] | None) -> list[int] | None:
if not repositories:
return None
rows = db.query(Repository).filter(Repository.name.in_(repositories)).all()
resolved_ids = [r.id for r in rows]
unresolved = set(repositories) - {r.name for r in rows}
if unresolved:
logger.warning("search.repositories_not_found", names=sorted(unresolved))
return resolved_ids or None
- 요청에는 레포 이름(repositories: ["auth-service"])이 담겨 오지만, Chroma 메타데이터에는 정수 ID가 저장되어 있음
- SQLite에서 이름을 ID로 변환하는 과정이 필요
- 존재하지 않는 레포 이름이 넘어오면 에러를 던지지 않고 경고 로그만 남김
- 나머지 유효한 레포에 대해서는 검색을 이어감
프론트엔드의 SSE 파싱
async function readSseStream(
body: ReadableStream<Uint8Array>,
onEvent: (event: SearchStreamEvent) => void
): Promise<void> {
const reader = body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const blocks = buffer.split("\n\n");
buffer = blocks.pop() ?? "";
for (const block of blocks) {
const event = parseSseBlock(block);
if (event) onEvent(event);
}
}
}
- TypeScript에서 SSE 스트림을 읽는 방법도 직접 구현
- EventSource 브라우저 API는 POST 요청을 지원하지 않아 fetch로 직접 처리
- 네트워크 청크 단위로 데이터가 도착하다 보면, SSE 이벤트 하나가 두 개의 네트워크 청크에 걸쳐 나뉘어 올 수 있음
\n\n으로 분리하되, 마지막 조각은 다음 청크와 합치기 위해 버퍼에 남겨둔다.
try {
return JSON.parse(data) as SearchStreamEvent;
} catch {
console.warn("[codemate-search] SSE 블록 JSON 파싱 실패, 건너뜀:", data.slice(0, 200));
return null;
}
- JSON 파싱 실패 시에는 해당 이벤트만 건너뛰고 스트림을 계속 유지
타임아웃 처리
async with asyncio.timeout(self.settings.request_timeout_seconds):
async for event in self._stream(request, repository_ids):
yield event
- LLM 응답 생성이 길어질 경우를 대비해 전체 스트리밍에 타임아웃 설정
- 타임아웃이 발생하면 error 이벤트를 전송하고 스트림을 종료
- 기본값은 환경변수 SEARCH_REQUEST_TIMEOUT_SECONDS로 주입
느낀 점
- SSE는 LLM 연동에 자연스러운 선택 — 응답 완성을 기다리지 않아도 되므로 UX가 크게 개선됨
- metadata를 먼저 보내는 순서가 중요 — 프론트가 참조 코드를 미리 렌더링 할 수 있고, 답변이 없어도 근거를 보여줄 수 있음
- EventSource를 못 쓰는 제약 — POST 요청이 필요한 경우 브라우저 내장 API 대신 fetch로 직접 SSE를 파싱해야 함
- Dense와 BM25의 필터 시점 차이는 구조적으로 자연스러운 결과 — 어느 쪽이 더 좋다기보다 각 저장소의 특성에 맞게 결정됨
'PROJECT > AI' 카테고리의 다른 글
| [GitHub App] GitHub App + Webhook으로 배운 자동 인덱싱 파이프라인 (0) | 2026.05.19 |
|---|---|
| [RAG] 사내 코드 검색 시스템을 만들며 배운 RAG 파이프라인 정리 (1) | 2026.05.15 |