"""채널 발견 → 확정 URL 크롤링 → fact·사진 후보 저장. 수집 결과는 사용자가 승인하기 전까지 사이트에 노출하지 않는다. """ import re import uuid from common import collect_diagnostics from common.database.db_session_manager import DB_SESSION_MNG from common.database.model.models import place_facts as facts_model from common.database.model.models import place_photos, place_channels, places, place_units from common.enums import DBWRType, ErrorType, LinkChannel, MediaStatus, PlaceCategory, PlaceStatus, SourceType from common.logger import LOG from common.utils.gtime import GTime from crud.fact_crud import FactCRUD from crud.place_crud import PlaceCRUD from common.category_schema import get_schema from services.collector import AdapterDisabled, AdapterNotFound, REGISTRY from services.collector import yanolja_adapter from services.external import naver_place_lookup, perplexity, tour_lookup from services.llm import provider from services.fact_service import FactService from router.v1.fact.protocol import Req_UpsertFact from common.job_errors import PermanentJobError _place_crud = PlaceCRUD() _fact_crud = FactCRUD() class CollectAborted(PermanentJobError): """재시도해도 소용없는 중단 — 잡의 last_error 로 남아 운영자가 본다.""" # ---- 단계 1: 채널 URL 발견 ------------------------------------------------- async def _add_link(place_id: str, channel, url: str, title: str, discovered_by, raw=None) -> bool: """링크 한 건 적재. 이미 있으면 False(유니크 충돌은 재수집의 정상 경로다).""" err = await DB_SESSION_MNG.execute_lambda_run( [place_channels.DBType()], [lambda s: _place_crud.add_link(s, place_channels( place_id=uuid.UUID(place_id), channel=channel.value, url=url, title=title, discovered_by=discovered_by.value, discovered_at=GTime.UTC(), raw=raw, ))], ) return err == ErrorType.SUCCESS async def discover_naver_place(place, place_id: str) -> str: """상호·주소로 네이버 플레이스 링크를 직접 찾아 등록한다. 반환값은 결과 상태다 — 호출측이 "자동으로 못 찾았다"를 사장님 화면까지 올려야 하므로 성공/실패를 bool 로 뭉개지 않는다: "resolved" 상호가 일치하는 place id 를 찾아 새로 등록·확정했다 "already" 같은 URL 이 이미 있었다(재수집의 정상 경로) — 확정만 다시 걸었다 "not_found" 상호가 일치하는 후보가 없다 → **자동 등록하지 않는다**(남의 가게 방지) ★ 왜 Perplexity 에 맡기지 않나 실측(2026-08-27 '도플로'·'버터브루'·'힐튼 가든 인 서울 강남'): Perplexity 는 야놀자만 물어오고 네이버 플레이스는 **한 건도** 못 찾았다. 도메인 필터를 넓혀도 안내 페이지 (pages.map.naver.com/useful-tips)가 걸릴 뿐이었다. 검색 언어모델은 "이 가게의 공식 플레이스 주소"를 안정적으로 집어내지 못한다 — 게다가 검색 1회당 요금이 붙는다. 그런데 우리는 이미 **이 가게가 누구인지 안다**(동일 업소 검증을 통과한 상호·주소). 추측할 이유가 없으므로 상호가 일치하는 place id 를 직접 해석한다. 일치하지 않으면 등록하지 않는다 — 남의 가게를 공식 채널로 붙이는 것이 여기서 가장 비싼 실수다. """ place_id_naver = await naver_place_lookup.find_place_id( place.name, place.road_address or place.address ) if not place_id_naver: # ★ 조용히 넘어가지 않는다 — 이 실패가 곧 "사장님이 주소를 직접 넣어야 한다"는 신호다. # 실측: '롯데호텔 서울' 이 여기서 떨어졌다(네이버 표기와 상호가 갈린다). LOG.i(f"[collect] 네이버 플레이스 자동 해석 실패 — 사장님이 주소를 붙여넣어야 한다 (place={place_id})") return "not_found" url = naver_place_lookup.place_url(place_id_naver) added = await _add_link( place_id, LinkChannel.NAVER_PLACE, url, f"{place.name} 네이버 플레이스", SourceType.API, ) # ★ 이 링크만은 자동 확정한다. # # 보통 확정은 사람이 한다 — 검색모델이 물어온 URL 은 "비슷한 이름의 다른 가게"일 수 있어서, # 그걸 자동 확정하면 남의 가게 페이지를 긁어 이 가게 사이트에 싣게 된다. # 그런데 이 링크는 그 경로로 온 것이 아니다. **상호가 정확히 일치할 때만** 등록되고 # (find_place_id 가 불일치면 None), 그 상호는 이미 동일 업소 검증을 통과한 값이다. # 근거가 사람 확인과 같은 수준이므로 클릭을 한 번 더 받는 것은 이득 없이 막기만 한다 — # 실제로 이 클릭 때문에 힐튼·도플로가 "수집했는데 0건"으로 끝났다. err, _rows = await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s: _place_crud.confirm_link_by_url(s, uuid.UUID(place_id), url, place.verified_by, GTime.UTC()), ) LOG.i(f"[collect] 네이버 플레이스 링크 {'등록·확정' if added else '확정'} — place {place_id_naver}") return "resolved" if added else "already" async def discover_official_site(place, place_id: str) -> str: """네이버 **지역검색**이 주는 업체 자체 홈페이지를 채널로 등록한다. 반환 규약은 `discover_naver_place` 와 같다("resolved" / "already" / "not_found"). ★ 왜 이게 필요한가 (실측 2026-09-10, 스테이,머뭄) 네이버 플레이스는 숙박에서 **fact 3건**(주차·와이파이·휠체어)밖에 주지 않는다. 체크인·체크아웃·취소규정·객실은 한 건도 없다. TourAPI 는 미등록 업소면 0건이고, 네이버 예약 페이지와 인스타그램은 robots 가 자동 수집을 금지한다. 그러면 남는 공개 출처는 **업소 자체 홈페이지 하나뿐**인데, 그 주소를 우리는 이미 받고 있었다 — 지역검색 응답의 `link` 다. 그걸 아무도 등록하지 않아 버려지고 있었다 (`external/naver.py` 머리주석: "채널 URL 발견에 쓸 수 있는 부수입이다"). 숙박은 자체 홈페이지 보유율이 3업종 중 가장 높다(표본 25건 중 19건). ★ 지어내지 않는다. 검색으로 URL 을 **추측**하는 Perplexity 경로와 다르다 — 이건 네이버가 그 업소 레코드에 달아 둔 값이고, 동일 업소 판정(`pick_match`)을 통과했을 때만 쓴다. 그래서 `discover_naver_place` 와 같은 근거로 자동 확정한다. """ from services.external import naver as naver_client client = naver_client.NaverLocalClient() if not client.enabled: return "not_found" address = place.road_address or place.address or "" # ★ `naver.region_key()` 를 쓰면 안 된다 — 그건 지명이 아니라 **행정구역 코드**('52군산시')다. # 검색어에 넣으면 후보가 0건이 된다(실측 2026-09-10). 사람이 읽는 지역 토막을 쓴다. query = " ".join(x for x in (place.name, naver_place_lookup.region_hint(address)) if x) try: candidates = await client.search_local(query) except Exception as ex: # noqa: BLE001 — 발견 실패가 수집을 죽이면 안 된다 collect_diagnostics.note_issue("naver_local_search", query, ex) return "not_found" match = naver_client.pick_match(place.name, candidates, address) if not match.is_matched or not match.place or not match.place.place_url: LOG.i(f"[collect] 자체 홈페이지 없음 — 네이버 레코드에 link 가 없거나 동일 업소 판정 실패 (place={place_id})") return "not_found" url = match.place.place_url channel = LinkChannel.INSTAGRAM if "instagram.com" in url.lower() else LinkChannel.OFFICIAL_SITE added = await _add_link(place_id, channel, url, f"{place.name} 공식 채널", SourceType.API, raw={"source": "naver_local"}) await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s: _place_crud.confirm_link_by_url(s, uuid.UUID(place_id), url, place.verified_by, GTime.UTC()), ) LOG.i(f"[collect] 자체 홈페이지 {'등록·확정' if added else '확정'} — {url}") return "resolved" if added else "already" async def discover_tour_api(place, place_id: str) -> str: """검증된 상호·좌표로 TourAPI 콘텐츠를 찾아 등록한다. 반환값은 discover_naver_place 와 같은 규약이다. ★ 왜 필요한가 (2026-08-31) `tour_api` 어댑터를 등록해도 **아무도 TourAPI URL 을 만들어주지 않아** 영영 안 불렸다. 채널 발견 프롬프트는 "야놀자·여기어때·네이버 플레이스" 만 찾기 때문이다. 실측: 롯데호텔 월드가 네이버 플레이스 하나로만 수집돼 fact 4건 · FAQ 2건에서 끝났다. 같은 업소를 TourAPI 로 받으면 **fact 93건 · 객실 11개 · 사진 30장**이다. ★ 자동 확정하는 이유는 네이버 플레이스와 같다 — 추측이 아니라 해석이기 때문이다. tour_lookup 이 **상호 일치 + 좌표 500m 이내**를 모두 통과할 때만 id 를 돌려준다 (부산 좌표로 '롯데호텔 월드' 를 조회하면 314km 떨어져 None 이 된다). 근거가 사람 확인과 같은 수준이라 클릭을 한 번 더 받는 것은 막기만 한다. """ if not tour_lookup.is_configured(): return "not_configured" try: found = await tour_lookup.find_content_id( place.name, PlaceCategory(place.category), latitude=place.latitude, longitude=place.longitude, ) except Exception as ex: # noqa: BLE001 — 발견 실패가 수집을 죽이면 안 된다 collect_diagnostics.note_issue("tour_api", place.name, ex) return "error" if not found: # 관광공사에 등록되지 않은 업소는 흔하다 — 실패가 아니라 정상적인 결말이다. LOG.i(f"[collect] TourAPI 미등록 업소 — 다른 채널로 진행 (place={place_id})") return "not_found" content_id, content_type = found url = f"tour://{content_type}/{content_id}" added = await _add_link( place_id, LinkChannel.ETC, url, f"{place.name} 한국관광공사 TourAPI", SourceType.API, ) err, _rows = await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s: _place_crud.confirm_link_by_url(s, uuid.UUID(place_id), url, place.verified_by, GTime.UTC()), ) LOG.i(f"[collect] TourAPI 링크 {'등록·확정' if added else '확정'} — contentId {content_id}") return "resolved" if added else "already" def _same_business(place_name: str, picked_name: str) -> bool: """상호 문자열이 같은 업소를 가리키는지 느슨하게 판정한다. ★ 야놀자는 상호로 정확히 질의하는 공식 API 가 없다 — 주소 검색 첫 결과가 진짜 이 업소인지 이 상호 비교로 한 번 더 확인한 뒤에만 자동 확정한다.""" def norm(s: str) -> str: return re.sub(r"[\s,.\-·]+", "", s or "").lower() p, q = norm(place_name), norm(picked_name) return bool(p) and bool(q) and (p in q or q in p) async def discover_yanolja(place, place_id: str) -> str: """주소로 야놀자(NOL) 상세페이지를 검색해 등록한다. 반환 규약은 discover_naver_place 와 같다. ★ 상호로 직접 질의하는 공식 API 가 없어 주소 검색 결과에 의존한다. 그래서 `discover_naver_place`·`discover_tour_api` 처럼 "해석"이라 부르기엔 근거가 약하다 — 검색 결과 상호가 place.name 과 겹치는지(`_same_business`) 확인했을 때만 자동 확정하고, 아니면 등록하지 않는다(남의 가게가 섞이는 것을 막는다). """ address = place.road_address or place.address if not address: return "not_found" try: found = await yanolja_adapter.search_by_address(address) except Exception as ex: # noqa: BLE001 — 발견 실패가 수집을 죽이면 안 된다 collect_diagnostics.note_issue("yanolja_search", place.name, ex) return "error" if not found: LOG.i(f"[collect] 야놀자 검색결과 없음 (place={place_id})") return "not_found" url, picked_name = found if not _same_business(place.name, picked_name): LOG.i(f"[collect] 야놀자 검색결과 상호 불일치 — 자동 등록 안 함 ({place.name!r} vs {picked_name!r})") return "not_found" added = await _add_link(place_id, LinkChannel.YANOLJA, url, f"{place.name} 야놀자", SourceType.API) await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s: _place_crud.confirm_link_by_url(s, uuid.UUID(place_id), url, place.verified_by, GTime.UTC()), ) LOG.i(f"[collect] 야놀자 링크 {'등록·확정' if added else '확정'} — {url}") return "resolved" if added else "already" async def discover_links(place, place_id: str, *, include_perplexity: bool = False) -> dict: """채널 URL 을 찾아 place_channels 에 적재한다(미확정 상태). ★ 순서가 곧 신뢰도다. 지금 자동 발견 경로는 **네이버 플레이스 직접 해석 하나뿐**이고, Perplexity 는 사용자가 추가 채널 탐색 옵션을 고른 회차에만 실행한다. 둘 다 실패하면 그건 정상적인 결말이다 — 사장님이 네이버 플레이스 주소를 붙여넣는 경로가 1순위이기 때문이다(화면 3단계가 그 입력을 맨 위에 둔다). ★ Perplexity 답변을 사실로 쓰지 않는다 — URL 발견 전용이다. 본문은 raw 에 박제만 한다.""" stat = {"discovered": 0, "skipped_duplicate": 0, "searches": 0, "enabled": False} # 네이버 플레이스는 검색모델에 맡기지 않고 직접 해석한다(위 주석 참고). # Perplexity 설정 여부와 무관하게 먼저 시도한다 — 키가 없어도 이건 된다. try: stat["naver_place"] = await discover_naver_place(place, place_id) if stat["naver_place"] == "resolved": stat["discovered"] += 1 except Exception as ex: # noqa: BLE001 — 발견 실패가 수집을 죽이면 안 된다 collect_diagnostics.note_issue("naver_place", place.name, ex) stat["naver_place"] = "error" # TourAPI 도 같은 성격의 '직접 해석' 이다 — 검색모델을 거치지 않고, 키가 있으면 항상 시도한다. # ★ 네이버보다 훨씬 많은 fact 를 준다(실측 롯데호텔 월드: 네이버 4건 vs TourAPI 93건). try: stat["tour_api"] = await discover_tour_api(place, place_id) if stat["tour_api"] == "resolved": stat["discovered"] += 1 except Exception as ex: # noqa: BLE001 collect_diagnostics.note_issue("tour_api", place.name, ex) stat["tour_api"] = "error" # 자체 홈페이지 — 네이버 지역검색이 이미 준 값이라 추가 요금이 없다(위 함수 머리주석). # ★ 숙박에서 이게 체크인·취소규정·객실을 가진 유일한 공개 출처인 경우가 많다. try: stat["official_site"] = await discover_official_site(place, place_id) if stat["official_site"] == "resolved": stat["discovered"] += 1 except Exception as ex: # noqa: BLE001 collect_diagnostics.note_issue("official_site", place.name, ex) stat["official_site"] = "error" # 야놀자(NOL) — 숙박 업종에서만 의미가 있고, 상호 대조 실패 시 등록하지 않는다(위 함수 참고). if PlaceCategory(place.category) == PlaceCategory.LODGING: try: stat["yanolja"] = await discover_yanolja(place, place_id) if stat["yanolja"] == "resolved": stat["discovered"] += 1 except Exception as ex: # noqa: BLE001 collect_diagnostics.note_issue("yanolja", place.name, ex) stat["yanolja"] = "error" # 오직 요청 옵션으로만 연다. 서버 env 로 일괄 활성화하면 일반 크롤링·재수집에서도 # 사용자가 모르는 유료 검색이 반복될 수 있으므로 COLLECT_USE_PERPLEXITY 는 더 쓰지 않는다. if not include_perplexity: LOG.i("[collect] 추가 채널 자동 발견 미선택 — Perplexity URL 검색 건너뜀") return stat stat["enabled"] = perplexity.is_configured() if not stat["enabled"]: LOG.i("[collect] Perplexity 미설정 — URL 발견 건너뜀(등록된 링크로 진행)") return stat try: found = await perplexity.discover_channels( place.name, address=place.road_address or place.address, category_hint=PlaceCategory(place.category).name, ) except perplexity.PerplexityNotConfigured: stat["enabled"] = False return stat except perplexity.PerplexityError as ex: # ★ 실패해도 파이프라인을 죽이지 않는다 — 이미 등록된 링크로 크롤링은 계속한다. collect_diagnostics.note_issue("perplexity_discover", place.name, ex) stat["error"] = str(ex)[:200] return stat stat["searches"] = found.search_count # 필터 탈락 내역 — 조용히 버리지 않는다. 운영자가 "왜 이 URL 이 빠졌나" 를 볼 수 있어야 한다. if hasattr(found, "reason_counts"): stat["filtered_out"] = found.reason_counts() now = GTime.UTC() for link in found.links: row = place_channels( place_id=uuid.UUID(place_id), channel=link.channel.value, url=link.url, title=link.title, discovered_by=SourceType.API.value, discovered_at=now, raw=found.raw, # 본문 + search_results 원문 — 환각 추적용 ) err = await DB_SESSION_MNG.execute_lambda_run( [place_channels.DBType()], [lambda s, r=row: _place_crud.add_link(s, r)], ) if err == ErrorType.SUCCESS: stat["discovered"] += 1 else: # (place_id, url) 부분 유니크 — 이미 있는 URL 이면 무시한다(재수집 시 정상 경로). stat["skipped_duplicate"] += 1 LOG.i(f"[collect] URL 발견 {stat['discovered']}건(중복 {stat['skipped_duplicate']}) · 검색 {stat['searches']}회") return stat # ---- 단계 2: 크롤링 대상 확정 ---------------------------------------------- async def confirm_targets(place, place_id: str, only_link_ids: list[str] | None) -> tuple[list, dict]: """크롤링할 링크를 고른다. ★ 사업장이 카카오 로컬로 동일 업소 검증을 통과해야만 여기까지 온다(호출측이 막는다). 검증된 사업장에 대해 도메인 필터를 통과한 URL 이므로 자동 확정한다. **그래도 여기서 나온 값은 전부 후보로 들어간다** — 사람이 승인해야 사이트에 나간다. 이 2중 방어가 '남의 가게 정보가 사이트에 실리는 것'을 막는 실제 장치다. """ stat = {"confirmed": 0, "already": 0, "unsupported": 0} err, links = await DB_SESSION_MNG.execute_lambda( place_channels.DBType(), DBWRType.DB_READ.value, lambda s: _place_crud.list_links(s, uuid.UUID(place_id), False), ) if err != ErrorType.SUCCESS: raise CollectAborted(f"채널 링크 조회 실패: {err.name}") wanted = set(only_link_ids or []) now = GTime.UTC() targets = [] for link in links: if wanted and str(link.link_id) not in wanted: continue # 처리할 어댑터가 없는 URL 은 확정하지 않는다 — 확정해봐야 긁지 못한다. if not REGISTRY.can_handle(link.url): stat["unsupported"] += 1 continue if link.confirmed_at is not None: stat["already"] += 1 targets.append(link) continue _e, rowcount = await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s, l=link: _place_crud.confirm_link(s, uuid.UUID(place_id), l.link_id, place.verified_by, now), ) if rowcount: stat["confirmed"] += 1 targets.append(link) LOG.i(f"[collect] 크롤링 대상 {len(targets)}건(신규확정 {stat['confirmed']} · 기확정 {stat['already']} " f"· 어댑터없음 {stat['unsupported']})") return targets, stat # ---- 단계 3: 크롤링 ------------------------------------------------------- async def coverage(place, place_id: str) -> dict: """이 사업장이 사이트를 만들 만큼 정보를 갖췄는지 — 업종 스키마의 필수 항목 기준. 이미 확보한 fact(노출값·후보 모두)의 key 를 세어 required 를 얼마나 덮었는지 본다.""" schema = get_schema(PlaceCategory(place.category)) required = set(schema.required_keys("place")) err, rows = await DB_SESSION_MNG.execute_lambda( facts_model.DBType(), DBWRType.DB_READ.value, lambda s: _fact_crud.list_facts(s, uuid.UUID(place_id), None, None, False, True), ) have = {r.key for r in (rows or []) if (r.value or "").strip()} if err == ErrorType.SUCCESS else set() covered = required & have return { "required": sorted(required), "missing": sorted(required - covered), "covered": len(covered), "total": len(required), "enough": not (required - covered), } async def fetch_one(link, category=None): """링크 하나를 긁는다. 실패해도 예외를 던지지 않는다 — 나머지 링크가 살아야 한다.""" try: adapter = REGISTRY.get_adapter(link.url) except (AdapterNotFound, AdapterDisabled) as ex: # 법무 검토 전이라 미등록된 어댑터(헤드리스 등)면 여기로 온다. LOG.w(f"[collect] 어댑터 없음 — 건너뜀 {link.url}: {type(ex).__name__}") return None, "no_adapter" try: source = await adapter.fetch(link.url, category) except Exception as ex: collect_diagnostics.note_issue("fetch", link.url, ex) return None, "failed" if not source.ok: collect_diagnostics.note_issue("fetch", link.url, RuntimeError(source.error)) return None, "failed" return source, "fetched" # ---- 단계 4: 하위 단위 시드 ------------------------------------------------ async def ensure_units(place_id: str, sources: list) -> dict: """수집된 unit_name(객실·메뉴·프로그램)을 place.units 에 맞춰 {이름: unit_id} 를 만든다. 수집 시점엔 unit_id 를 모르므로 이름으로 매핑한다. 없는 이름만 새로 만든다.""" names = [] for source in sources: for name in source.unit_names(): if name not in names: names.append(name) if not names: return {} err, rows = await DB_SESSION_MNG.execute_lambda( place_units.DBType(), DBWRType.DB_READ.value, lambda s: _place_crud.list_units(s, uuid.UUID(place_id)), ) existing = {r.name: r.unit_id for r in (rows or [])} if err == ErrorType.SUCCESS else {} created = 0 for order, name in enumerate(names): if name in existing: continue row = place_units(place_id=uuid.UUID(place_id), name=name, sort_order=order) run_err = await DB_SESSION_MNG.execute_lambda_run( [place_units.DBType()], [lambda s, r=row: _place_crud.add_unit(s, r)], ) if run_err == ErrorType.SUCCESS: existing[name] = row.unit_id created += 1 if created: LOG.i(f"[collect] 하위 단위 {created}건 생성") return existing # ---- 단계 5: fact / 사진 적재 ---------------------------------------------- async def store_facts(user_info, place_id: str, sources: list, unit_map: dict) -> dict: """수집값을 fact 로 적재한다. ★ FactService 를 그대로 통과시킨다 — 업종 스키마 검증 · 출처 필수 · LLM 제한 · 정정본 보호 · 노출값 유지가 전부 거기 있다. 크롤러가 우회할 수 있는 뒷문을 만들지 않는다. """ service = FactService(_fact_crud, _place_crud) stat = {"stored": 0, "refreshed": 0, "candidate": 0, "rejected": 0, "by_reason": {}} for source in sources: for cf in source.facts: unit_id = unit_map.get(cf.unit_name) if cf.scope == "unit" else None if cf.scope == "unit" and unit_id is None: stat["rejected"] += 1 stat["by_reason"]["UNIT_NOT_FOUND"] = stat["by_reason"].get("UNIT_NOT_FOUND", 0) + 1 continue req = Req_UpsertFact( key=cf.key, value=cf.value, unit_id=unit_id, source_type=SourceType.CRAWL, source_url=cf.source_url, ) res = await service.upsert_fact(user_info, place_id, req) if not res.result.success: stat["rejected"] += 1 stat["by_reason"][res.result.desc] = stat["by_reason"].get(res.result.desc, 0) + 1 continue stat["stored"] += 1 if res.outcome is not None and res.outcome.name == "REFRESHED": stat["refreshed"] += 1 elif res.outcome is not None and res.outcome.name.startswith("CANDIDATE"): stat["candidate"] += 1 return stat async def store_media(place_id: str, sources: list, unit_map: dict) -> dict: """수집한 사진을 적재한다. 같은 origin_url 은 다시 넣지 않는다(재수집 멱등). ★ source_type=CRAWL 과 origin_url 을 반드시 남긴다 — 이미지 재게시 권리가 미결이라 결론에 따라 발행 시 통째로 걸러낼 수 있어야 한다(docs/DECISIONS.md 1-2). ★ Vision 분석 전이므로 status 는 PENDING_REVIEW 다. 사람 확인 큐로 간다.""" from sqlalchemy import select err, rows = await DB_SESSION_MNG.execute_lambda( place_photos.DBType(), DBWRType.DB_READ.value, lambda s: DB_SESSION_MNG.execute( s, select(place_photos).where(place_photos.place_id == uuid.UUID(place_id), place_photos.deleted == False) # noqa: E712 ), ) seen = {r.origin_url for r in (rows or []) if r.origin_url} if err == ErrorType.SUCCESS else set() stat = {"stored": 0, "skipped_duplicate": 0} for source in sources: for order, cm in enumerate(source.media): if cm.origin_url in seen: stat["skipped_duplicate"] += 1 continue row = place_photos( place_id=uuid.UUID(place_id), unit_id=unit_map.get(cm.unit_name), url=cm.origin_url, # Vision·재게시 결론 전까지는 원본 URL 을 그대로 둔다 origin_url=cm.origin_url, source_type=SourceType.CRAWL.value, source_url=cm.source_url, label=cm.label, status=MediaStatus.PENDING_REVIEW.value, sort_order=order, ) run_err = await DB_SESSION_MNG.execute_lambda_run( [place_photos.DBType()], [lambda s, r=row: DB_SESSION_MNG.insert(s, r, raise_error=False)], ) if run_err == ErrorType.SUCCESS: seen.add(cm.origin_url) stat["stored"] += 1 return stat # ---- 오케스트레이션 -------------------------------------------------------- async def run_collect(job: dict) -> dict: """COLLECT 잡 핸들러. 반환값이 jobs.result 에 저장돼 폴링·감사에 쓰인다.""" with collect_diagnostics.collecting(): result = await _run_collect(job) issues = collect_diagnostics.snapshot() if issues: result["issues"] = issues return result async def _run_collect(job: dict) -> dict: payload = job["payload"] place_id = payload["place_id"] owner_user_id = payload["owner_user_id"] err, place = await DB_SESSION_MNG.execute_lambda( places.DBType(), DBWRType.DB_READ.value, lambda s: _place_crud.get_place(s, uuid.UUID(owner_user_id), uuid.UUID(place_id)), ) if err != ErrorType.SUCCESS or place is None: raise CollectAborted(f"사업장을 찾을 수 없다: {place_id}") # ★ 동일 업소 검증 게이트 — 잡 실행 시점에도 다시 확인한다(적재 후 취소됐을 수 있다). if place.verified_at is None: raise CollectAborted("동일 업소 검증(verify) 전에는 수집하지 않는다 — 남의 가게가 섞인다") result: dict = {"place_id": place_id} result["discover"] = await discover_links( place, place_id, include_perplexity=bool(payload.get("discover_channels", False)), ) targets, result["confirm"] = await confirm_targets(place, place_id, payload.get("link_ids")) # ★ 자동 해석 실패를 사장님에게 알리는 유일한 창구. # 지금 네이버 플레이스는 "URL 만 있으면 100%, 상호로 찾는 건 4곳 중 3곳" 이다. # 못 찾았을 때 조용히 넘어가면 사장님은 "수집했는데 아무것도 안 나왔다"만 본다 — # 실제로 해야 할 일(네이버 지도에서 내 가게 주소를 복사해 붙여넣기)을 화면이 말해줄 수 # 있도록 결과에 싣는다. 판정 기준은 "지금 긁을 네이버 플레이스 링크가 있느냐" 다: # 자동 해석이 실패해도 사장님이 이미 붙여넣었으면 알릴 이유가 없다. result["naver_place_missing"] = not any( link.channel == LinkChannel.NAVER_PLACE.value for link in targets ) if not targets: result["note"] = "크롤링 대상이 없다(어댑터가 처리할 수 있는 확정 URL 0건)" await _finish(place_id, owner_user_id, PlaceStatus.REVIEW) return result # ★ 이미 충분하면 크롤링 자체를 건너뛴다(force 가 아닐 때). before = await coverage(place, place_id) if before["enough"] and not payload.get("force"): result["coverage"] = before result["note"] = "이미 필수 항목이 다 차 있다 — 크롤링 생략(force=true 로 강제 가능)" await _finish(place_id, owner_user_id, PlaceStatus.REVIEW) LOG.i(f"[collect] 충분함 place={place_id} {before['covered']}/{before['total']} — 크롤링 생략") return result from common.models.gmodel import UserInfo actor = UserInfo( # ★ 잡이 쓰는 신원. user_id 는 **사업장 주인**이어야 한다 — FactService 가 이 값으로 # 사업장을 스코프하고(fact_service._load_place) verified_by 에도 그대로 박는다. # 회사를 걷어내기 전에는 스코프가 company_id 였고 여기엔 요청자·검증자·랜덤 uuid 가 # 순서대로 들어갔다. 그 랜덤 uuid 가 이제는 "남의 사업장" 이 되어 조회가 0건이 된다. user_id=owner_user_id, id="collector", role=1, ) # ★ 링크를 하나씩 긁고 **바로 적재한 뒤** 충분한지 본다. # 전부 긁어놓고 나중에 적재하면 "이미 충분한데 더 긁는" 낭비를 못 막는다. # 단계마다 커밋되므로 중간에 죽어도 여기까지의 결과는 남는다. fetch_stat = {"fetched": 0, "failed": 0, "no_adapter": 0, "skipped_enough": 0, "stopped_early": False} facts_stat = {"stored": 0, "refreshed": 0, "candidate": 0, "rejected": 0, "by_reason": {}} media_stat = {"stored": 0, "skipped_duplicate": 0} unit_total = 0 for index, link in enumerate(targets): if index > 0: cov = await coverage(place, place_id) if cov["enough"]: fetch_stat["skipped_enough"] = len(targets) - index fetch_stat["stopped_early"] = True LOG.i(f"[collect] 필수 항목 {cov['covered']}/{cov['total']} 확보 — " f"남은 링크 {fetch_stat['skipped_enough']}건 크롤링 생략") break source, outcome = await fetch_one(link, PlaceCategory(place.category)) fetch_stat[outcome] += 1 if source is None: continue # ★ 수집 원문을 링크에 남긴다 — fact 가 아니라 '생성 근거' 자리다. # intro 같은 allow_llm 필드에 원문을 넣으면 발행본이 원문으로 덮인다(2026-08-31 사고). if (source.text or "").strip(): await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda sess, u=link.url, t=source.text: _place_crud.set_link_raw( sess, uuid.UUID(place_id), u, {"text": t[:8000]}, ), ) await _store_booking_link(place, place_id, source) unit_map = await ensure_units(place_id, [source]) unit_total = max(unit_total, len(unit_map)) f = await store_facts(actor, place_id, [source], unit_map) m = await store_media(place_id, [source], unit_map) for k in ("stored", "refreshed", "candidate", "rejected"): facts_stat[k] += f[k] for reason, n in f["by_reason"].items(): facts_stat["by_reason"][reason] = facts_stat["by_reason"].get(reason, 0) + n media_stat["stored"] += m["stored"] media_stat["skipped_duplicate"] += m["skipped_duplicate"] result["fetch"] = fetch_stat result["units"] = unit_total result["facts"] = facts_stat result["media"] = media_stat result["coverage"] = await coverage(place, place_id) # 사진이 들어왔으면 분석을 이어서 건다 — 수집과 분석은 각각 몇 분이라 한 잡에 묶지 않는다. # (묶으면 분석에서 죽었을 때 수집까지 다시 하게 되고, 유료 API 를 두 번 태운다.) if result["media"]["stored"] > 0: result["vision_job_id"] = await _enqueue_vision(place_id, owner_user_id) # ★ 지역 데이터(주변 맛집·관광지·축제 + 지역 이야기)를 **여기서** 건다. # 수집이 끝난 시점이 좌표·행정구역이 확정되는 가장 이른 자리다. 사장님이 템플릿을 고르는 # 동안(Step4) 백그라운드로 돌아, 생성 단계(Step5)에 닿을 즈음이면 대개 끝나 있다 — # 전에는 에디터에 들어간 뒤에야 시작해서 첫 화면이 늘 절반만 그려졌다. from services import story_service result["local_job_id"] = await story_service.enqueue_region_job(place) # ── 업소 조사 — 소개문을 쓸 재료 ────────────────────────────────── # ★ 수집이 끝난 **뒤**에 한다. 앞에서 하면 네이버·TourAPI 가 이미 준 것을 다시 묻는 # 꼴이고, 검색 요금이 그만큼 헛돈다. 수집이 얇게 끝났을 때 그 구멍을 메우는 자리다. # ★ fact 를 만들지 않는다(place_research 머리주석) — 소개문 생성의 근거만 쌓는다. # ★ 실패해도 수집은 성공이다. 재료가 적을 뿐 발행은 된다. try: from services import place_research result["research"] = await place_research.research_place(place, place_id) except Exception as ex: # noqa: BLE001 collect_diagnostics.note_issue("place_research", place.name, ex) result["research"] = {"error": f"{type(ex).__name__}: {ex}"} await _finish(place_id, owner_user_id, PlaceStatus.REVIEW) LOG.i(f"[collect] 완료 place={place_id} fact {result['facts']['stored']}건 · 사진 {result['media']['stored']}장") return result async def _store_booking_link(place, place_id: str, source) -> bool: """수집 중 채널이 알려준 예약 주소를 **예약 채널 링크**로 남긴다. ★ 왜 필요한가 (실측 2026-09-08) 발행본의 "예약" 버튼이 네이버 플레이스 링크를 그대로 열었다. 그 링크는 잘해야 플레이스 홈이라 예약까지 한 번 더 눌러야 하고, 자동 발견이 검색 URL(`map.naver.com/p/search/…`)을 물어온 경우에는 **검색 결과 화면**이 뜬다. 예약하려고 누른 손님이 검색 결과를 만나면 거기서 끝난다. ★ 주소를 만들지 않는다. 플레이스 응답의 `naverBookingUrl` 을 그대로 쓴다 (naver_place_adapter._booking_url 머리주석). 예약을 받지 않는 업소에는 이 값이 없고, 없으면 링크도 없다 — 없는 예약 창구를 만들어내지 않는다. ★ 자동 확정한다. 근거는 `discover_naver_place` 와 같다 — 이 URL 은 **이미 확정된** 플레이스 페이지가 자기 예약 주소로 내놓은 값이라, 남의 가게가 섞일 경로가 없다. 여기서 클릭을 한 번 더 받으면 사장님이 확정을 안 한 사이트는 예약 버튼이 계속 검색 화면으로 간다. """ url = (getattr(source, "booking_url", None) or "").strip() if not url: return False added = await _add_link( place_id, LinkChannel.NAVER_BOOKING, url, f"{place.name} 네이버 예약", SourceType.CRAWL, ) await DB_SESSION_MNG.execute_lambda_claim( place_channels.DBType(), lambda s: _place_crud.confirm_link_by_url(s, uuid.UUID(place_id), url, place.verified_by, GTime.UTC()), ) LOG.i(f"[collect] 네이버 예약 링크 {'등록·확정' if added else '확정'} — {url}") return added async def _finish(place_id: str, owner_user_id: str, status: PlaceStatus): """수집이 끝나면 사업장을 검수 대기로 돌린다 — 수집값은 전부 후보라 사람이 봐야 한다.""" await DB_SESSION_MNG.execute_lambda_claim( places.DBType(), lambda s: _place_crud.update_place( s, uuid.UUID(owner_user_id), uuid.UUID(place_id), {"status": status.value} ), ) async def _enqueue_vision(place_id: str, owner_user_id: str) -> str | None: """사진 분석 잡을 적재한다. 키가 없거나 중복이면 조용히 건너뛴다(수집 자체는 성공이다).""" from common.enums import JobType from crud.job_crud import JobQueue from services.external import gemini from services.job_service import enqueue_job if not gemini.is_configured(): LOG.i(f"[collect] {provider.missing_key()} 미설정 — 사진 분석 건너뜀(사진은 확인 큐에 남는다)") return None job_id, _created = await enqueue_job( JobQueue(), JobType.VISION, {"place_id": place_id, "owner_user_id": owner_user_id}, dedupe_key=f"vision:{place_id}", ) return job_id