"""잡 핸들러 — job_type 별 처리. 현재는 SEARCH(검색)만. 검색 핸들러 = 한정된 재정제 루프 + 명시적 outcome: 0. 네거티브 캐시 확인(최근 not_found면 즉시 반환) 각 라운드(원본 → 정밀(LLM) → 광역, 최대 max_rounds): 소스 동시 검색 → 필터 → 이상치 → AI 같은상품 판정 ├ 매칭 있음 → DONE(outcome=found), 조기 종료 ├ 0매칭 + 기술적 실패(차단/예외) 있음 → raise → 큐가 잡 전체 백오프 재시도(→소진 시 DEAD) └ 0매칭 + 소스 정상 → 다음 라운드 라운드 소진 → DONE(outcome=not_found) + 네거티브 캐시 기록 두 재시도 축을 분리한다: 기술적(큐 attempts/백오프) ≠ 검색어(refine 라운드, 유한). '못 찾음'은 정상 종료(DONE)지 dead-letter 가 아니다. """ import asyncio from common.enums import JobType from common.logger import LOG from services.search.contract import SearchAdapter from services.search.util import parse_price from services.pipeline.core import apply_filters, rank_result def build_search_handler( adapters: dict[str, SearchAdapter], sources: list[str] | None = None, limit: int = 40, top_n: int = 5, judge=None, keyword_gen=None, max_rounds: int = 3, neg_cache=None, ): """검색 핸들러 생성. judge: SimilarityJudge(같은 상품 판정) / keyword_gen: KeywordGenerator(정밀·광역 재검색어) / neg_cache: NegativeCache(TTL not_found 캐시). 모두 선택 — 없으면 해당 단계 생략.""" use = list(sources) if sources else list(adapters.keys()) async def _search_round(query: str): """한 라운드: 모든 소스 동시 검색 → (products, per_source, tech_failed).""" results = await asyncio.gather(*[adapters[s].search(query, limit=limit) for s in use], return_exceptions=True) products, per_source, tech_failed = [], {}, False for src, res in zip(use, results): if isinstance(res, Exception): tech_failed = True per_source[src] = {"error": f"{type(res).__name__}: {res}"} LOG.w(f"[{src}] 검색 실패: {type(res).__name__}: {res}") else: products.extend(res) per_source[src] = {"count": len(res)} return products, per_source, tech_failed async def _round_queries(base_query: str, target: dict): """라운드 쿼리 지연 생성: 원본 → (0매칭 시에만 LLM 호출로) 정밀 → 광역.""" yield ("original", base_query) if keyword_gen is not None: kw = await keyword_gen.generate(target) # 원본이 실패해 여기까지 온 경우에만 호출됨 seen = {base_query} for label, q in (("precise", kw.precise), ("broad", kw.broad)): q = (q or "").strip() if q and q not in seen: seen.add(q) yield (label, q) async def handler(job: dict) -> dict: if job["job_type"] != JobType.SEARCH.value: raise ValueError(f"unsupported job_type: {job['job_type']}") payload = job.get("payload") or {} base_query = (payload.get("product_name") or "").strip() if not base_query: raise ValueError("empty product_name") target = {k: payload.get(k, "") for k in ("product_name", "model", "specification", "company")} base_price = parse_price(payload.get("price")) cache_key = payload.get("product_code") or base_query # 0) 네거티브 캐시 — 최근 not_found면 재검색 생략 if neg_cache is not None and await neg_cache.is_negative(cache_key): return {"outcome": "not_found", "cached": True, "query": base_query, "rounds_tried": 0, "lowest": None, "top": [], "stages": [], "sources": {}} rounds_done = 0 last_stages, last_sources = [], {} async for label, query in _round_queries(base_query, target): if rounds_done >= max_rounds: break rounds_done += 1 products, per_source, tech_failed = await _search_round(query) candidates, stages = apply_filters(products, base_price=base_price) if judge is not None and candidates: verdicts = await judge.judge(target, candidates) matched = [c for c, v in zip(candidates, verdicts) if v.is_match] stages.append({"stage": "ai_match", "in": len(candidates), "out": len(matched)}) candidates = matched last_stages, last_sources = stages, per_source if candidates: # 찾음 → 조기 종료 result = rank_result(candidates, len(products), stages, top_n) result.update(outcome="found", query=query, round=label, rounds_tried=rounds_done, sources=per_source) return result if tech_failed: # 0매칭인데 소스가 죽어 있었음 → '없음'이라 단정 불가 → 기술 재시도 raise RuntimeError(f"기술적 실패로 0매칭(round={label}) — 잡 재시도: {per_source}") # 모든 라운드 클린 0매칭 → 정상 not_found 종료 if neg_cache is not None: await neg_cache.put(cache_key, reason=f"not_found after {rounds_done} rounds") result = rank_result([], 0, last_stages, top_n) result.update(outcome="not_found", query=base_query, rounds_tried=rounds_done, sources=last_sources) return result return handler