"""잡 핸들러 — 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 import time from common.enums import JobType from common.logger import LOG from services.metrics import SearchMetrics from services.search.contract import SearchAdapter, NormalizedProduct from services.search.card_parser import canonical_mall, MALL_BY_SOURCE from services.search.util import parse_price from services.pipeline.core import apply_filters, rank_result, summarize_by_mall # 데드라인 초과로 버린 폴백 태스크의 강한 참조(asyncio 는 태스크를 약참조만 유지 — 없으면 GC 로 중도 파괴될 수 있음). # 완료 시 콜백에서 스스로 제거된다. 테스트는 이 집합을 gather 해 잔여 태스크를 배수(drain)할 수 있다. _abandoned_fallbacks: set[asyncio.Task] = set() def _reap_abandoned(task: asyncio.Task): """버려진 폴백 태스크 종료 시 예외를 회수 — 'Future exception was never retrieved' 노이즈 방지.""" _abandoned_fallbacks.discard(task) if task.cancelled(): return ex = task.exception() if ex is not None: LOG.d(f"[fallback] 데드라인 초과 태스크 종료(무시): {type(ex).__name__}") def _price_snapshot(matched: list[NormalizedProduct]) -> dict: """매칭 목록에서 소스별 최저가 + 전체 최저가 스냅샷을 만든다(price_history 기록용).""" def lowest(src): items = [p for p in matched if p.source == src] return min(items, key=lambda p: p.price) if items else None n, c = lowest("naver"), lowest("coupang") # 최종 최저가는 소스 무관 전체 매칭 중 최저(G마켓·옥션·11번가 등 폴백 포함). f = min(matched, key=lambda p: p.price) if matched else None return { "matched_count": len(matched), "naver_lowest": n.price if n else None, "naver_name": n.name if n else None, "naver_url": n.detail_url if n else None, "coupang_lowest": c.price if c else None, "coupang_name": c.name if c else None, "coupang_url": c.detail_url if c else None, "final_lowest": f.price if f else None, "final_source": f.source if f else None, "by_mall": summarize_by_mall(matched), # 몰별 최저가 스냅샷(열린 스키마) } 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, history=None, fallback_adapters: dict[str, SearchAdapter] | None = None, ai_model: str = "", proxy_cost_per_gb: float = 0.0, fallback_deadline_sec: float = 15.0, ): """검색 핸들러 생성. judge: SimilarityJudge(같은 상품 판정) / keyword_gen: KeywordGenerator(정밀·광역 재검색어) / neg_cache: NegativeCache(TTL not_found 캐시) / history: PriceHistory(최저가 스냅샷) / fallback_adapters: 오픈마켓 크롤(gmarket/auction/st11) — 네이버가 그 몰을 커버 못 했을 때만 크롤(폴백). 모두 선택 — 없으면 해당 단계 생략.""" use = list(sources) if sources else list(adapters.keys()) fallbacks = fallback_adapters or {} async def _record_history(product_code: str, job_id, outcome: str, matched: list): if history is None: return event = {"product_code": product_code, "job_id": job_id, "outcome": outcome, **_price_snapshot(matched)} try: await history.record(event) except Exception as ex: LOG.e_no_callstack(f"[history] 스냅샷 기록 실패(무시): {ex}") async def _timed_search(adapter, query: str, source: str, metrics: SearchMetrics, crawl: bool): """어댑터 검색 1건을 타이밍+바이트 계측하며 실행. 예외는 그대로 전파(호출부가 처리).""" via_proxy = bool(getattr(adapter, "uses_proxy", False)) t0 = time.monotonic() try: res = await adapter.search(query, limit=limit) metrics.add_fetch(source, getattr(adapter, "last_bytes", 0), int((time.monotonic() - t0) * 1000), crawl=crawl, via_proxy=via_proxy) return res except BaseException: # CancelledError(데드라인 취소) 포함 — 소요/바이트는 계측하고 재전파 metrics.add_fetch(source, getattr(adapter, "last_bytes", 0), int((time.monotonic() - t0) * 1000), crawl=crawl, via_proxy=via_proxy) raise async def _search_round(query: str, metrics: SearchMetrics): """한 라운드: 모든 소스 동시 검색 → (products, per_source, tech_failed). 소스별 시간/바이트 계측.""" results = await asyncio.gather(*[_timed_search(adapters[s], query, s, metrics, False) 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 _match(target: dict, products: list, base_price, metrics: SearchMetrics): """필터 → (있으면) AI 같은상품 판정 → 매칭 후보. AI 토큰은 metrics 에 누적. ⚠️ 동시성 불변식: judge.judge 가 last_usage 를 세팅한 뒤 여기서 읽기까지 await 이 없어야 한다(asyncio 협조 스케줄링상 그 사이 다른 코루틴이 last_usage 를 덮어쓸 수 없음). 병렬 폴백 안전.""" candidates, _ = apply_filters(products, base_price=base_price) if judge is not None and candidates: verdicts = await judge.judge(target, candidates) metrics.add_ai(judge.last_usage) # ← await 직후 즉시 읽음(사이에 await 금지) candidates = [c for c, v in zip(candidates, verdicts) if v.is_match] return candidates async def _enrich_with_fallback(target: dict, query: str, matched: list, base_price, metrics: SearchMetrics): """네이버가 커버 못 한 오픈마켓만 직접 크롤(폴백) → 같은상품 판정 후 병합. 사용자 규칙: '네이버로 그 몰 값 확보 성공 → 그 값, 실패(몰 없음) → 실사이트 크롤'. 미커버 몰들을 **동시 크롤**한다(각 어댑터=별 브라우저 인스턴스라 병렬 안전, 소요=합→최댓값).""" if not fallbacks: return matched covered = {canonical_mall(p) for p in matched} todo = [(src, ad) for src, ad in fallbacks.items() if MALL_BY_SOURCE.get(src, src) not in covered] if not todo: return matched async def _crawl_match(src, adapter): # 폴백은 '있으면 좋은' 보강이라 데드라인을 건다 — 초과 시 그 몰만 스킵(전체 지연에 상한). # cancel 하지 않고 버린다(asyncio.wait): in-flight page.goto 를 취소하면 patchright 내부 # future 가 미회수 예외 노이즈를 남기고 페이지가 어중간한 상태로 남는다. 버려진 크롤은 # 백그라운드에서 자체 타임아웃(goto 40s 등)으로 끝나고 _reap_abandoned 가 예외를 회수한다. task = asyncio.ensure_future(_timed_search(adapter, query, src, metrics, crawl=True)) done, _ = await asyncio.wait({task}, timeout=fallback_deadline_sec) if not done: _abandoned_fallbacks.add(task) task.add_done_callback(_reap_abandoned) LOG.w(f"[fallback:{src}] 데드라인 {fallback_deadline_sec:.0f}s 초과 → 스킵(크롤은 백그라운드 종료)") return [] try: crawled = task.result() except Exception as ex: LOG.w(f"[fallback:{src}] 크롤 실패(무시): {type(ex).__name__}: {ex}") return [] hits = await _match(target, crawled, base_price, metrics) if hits: LOG.d(f"[fallback:{src}] 크롤 {len(crawled)}건 중 같은상품 {len(hits)}건 병합") return hits results = await asyncio.gather(*[_crawl_match(s, a) for s, a in todo]) for hits in results: matched = matched + hits return matched async def _round_queries(base_query: str, target: dict, metrics: SearchMetrics): """라운드 쿼리 지연 생성: 원본 → (0매칭 시에만 LLM 호출로) 정밀 → 광역.""" yield ("original", base_query) if keyword_gen is not None: kw = await keyword_gen.generate(target) # 원본이 실패해 여기까지 온 경우에만 호출됨 metrics.add_ai(keyword_gen.last_usage) 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 metrics = SearchMetrics(ai_model, proxy_cost_per_gb) # 검색 1건의 리소스/비용/시간 계측 # 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": {}, "metrics": metrics.snapshot()} rounds_done = 0 last_stages, last_sources = [], {} async for label, query in _round_queries(base_query, target, metrics): if rounds_done >= max_rounds: break rounds_done += 1 products, per_source, tech_failed = await _search_round(query, metrics) candidates, stages = apply_filters(products, base_price=base_price) if judge is not None and candidates: verdicts = await judge.judge(target, candidates) metrics.add_ai(judge.last_usage) 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: # 찾음 → 오픈마켓 폴백 보강 후 종료 before = len(candidates) candidates = await _enrich_with_fallback(target, query, candidates, base_price, metrics) if len(candidates) > before: stages.append({"stage": "fallback_crawl", "in": before, "out": len(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, metrics=metrics.snapshot()) await _record_history(cache_key, job.get("job_id"), "found", candidates) 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, metrics=metrics.snapshot()) await _record_history(cache_key, job.get("job_id"), "not_found", []) return result return handler