본문 바로가기
PROJECT/AI

[Search API] SSE 스트리밍으로 설계한 코드 검색 API와 레포 필터 구조

by 놀고 쉬고 싶은 개발자 2026. 5. 22.

이전 글에서 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의 필터 시점 차이는 구조적으로 자연스러운 결과 — 어느 쪽이 더 좋다기보다 각 저장소의 특성에 맞게 결정됨