컨테이너를 재생성하면 docker logs 가 사라져 크롤링이 잘 됐는지 확인할 방법이 없었다. 실측(09-30 버터브루): 수집 도중 배포로 API 가 80초 끊기자 화면이 사진 분석 폴링을 3회 실패 후 포기해 사진 10장이 DB 에 있는데도 0장으로 보였다. - services/activity_feed.py: 기존 alert_outbox·장애 채널로 kind=activity 이벤트 적재(중복 억제 없음) - collect_service: 수집 시작 · 채널 크롤링 실패 즉시 · 완료 요약(채널별·fact·사진·누락 필수항목·소요초) · 실패 - vision_service: 사진 분석 결과 - build_service · rollback_service: 첫 발행/재발행/빌드만/되돌리기 시작·끝 — URL · 굽기 사진 미러링 수 - site/prerender.ts: 렌더 보고서에 사진 미러링 수(media) 추가 - frontend pollJob: 연속 3회 실패여도 3분간 무응답일 때만 unreachable - frontend collectJobs: 사진 분석 폴링을 못 끝내도 사진 목록을 다시 읽고 경고 - test_build_publish: 게이트 반려 검사에서 activity 이벤트는 제외 관련 테스트 17개 파일 295 passed · 2 failed(test_search_console_service — 변경 전에도 실패) frontend tsc·eslint, site tsc·eslint 통과 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
691 lines
30 KiB
Python
691 lines
30 KiB
Python
"""채널 발견 → 확정 URL 크롤링 → fact·사진 후보 저장."""
|
|
import re
|
|
import time
|
|
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 import activity_feed
|
|
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:
|
|
"""링크 한 건 적재."""
|
|
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:
|
|
"""상호·주소로 네이버 플레이스 링크를 직접 찾아 등록한다."""
|
|
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,
|
|
)
|
|
|
|
# 이 링크만은 자동 확정한다.
|
|
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:
|
|
"""네이버 **지역검색**이 주는 업체 자체 홈페이지를 채널로 등록한다."""
|
|
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군산시')다.
|
|
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 콘텐츠를 찾아 등록한다."""
|
|
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:
|
|
"""상호 문자열이 같은 업소를 가리키는지 느슨하게 판정한다."""
|
|
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) 상세페이지를 검색해 등록한다."""
|
|
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 에 적재한다(미확정 상태)."""
|
|
stat = {"discovered": 0, "skipped_duplicate": 0, "searches": 0, "enabled": False}
|
|
|
|
# 네이버 플레이스는 검색모델에 맡기지 않고 직접 해석한다(위 주석 참고).
|
|
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 도 같은 성격의 '직접 해석' 이다 — 검색모델을 거치지 않고, 키가 있으면 항상 시도한다.
|
|
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"
|
|
|
|
# 오직 요청 옵션으로만 연다.
|
|
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
|
|
# 필터 탈락 내역 — 조용히 버리지 않는다.
|
|
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]:
|
|
"""크롤링할 링크를 고른다."""
|
|
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:
|
|
"""이 사업장이 사이트를 만들 만큼 정보를 갖췄는지 — 업종 스키마의 필수 항목 기준."""
|
|
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} 를 만든다."""
|
|
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 로 적재한다."""
|
|
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:
|
|
"""수집한 사진을 적재한다."""
|
|
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 잡 핸들러."""
|
|
started = time.monotonic()
|
|
with collect_diagnostics.collecting():
|
|
try:
|
|
result = await _run_collect(job)
|
|
except Exception as ex:
|
|
await activity_feed.post(
|
|
"❌ 수집 실패",
|
|
[f"place_id={job['payload'].get('place_id')}", f"{type(ex).__name__}: {ex}",
|
|
f"{time.monotonic() - started:.0f}초"],
|
|
)
|
|
raise
|
|
issues = collect_diagnostics.snapshot()
|
|
if issues:
|
|
result["issues"] = issues
|
|
await _post_collect_summary(result, issues, time.monotonic() - started)
|
|
return result
|
|
|
|
|
|
async def _post_collect_summary(result: dict, issues: list[dict], elapsed: float) -> None:
|
|
lines = [result.get("label") or f"place_id={result.get('place_id')}"]
|
|
if result.get("note"):
|
|
lines.append(result["note"])
|
|
for row in result.get("channels") or []:
|
|
lines.append(row)
|
|
facts = result.get("facts")
|
|
if facts:
|
|
reasons = ", ".join(f"{k} {v}" for k, v in facts["by_reason"].items())
|
|
lines.append(f"fact 저장 {facts['stored']} · 반려 {facts['rejected']}" + (f" ({reasons})" if reasons else ""))
|
|
media = result.get("media")
|
|
if media:
|
|
lines.append(f"사진 {media['stored']}장 (중복 {media['skipped_duplicate']})")
|
|
cov = result.get("coverage")
|
|
if cov and cov["missing"]:
|
|
lines.append(f"필수 항목 {cov['covered']}/{cov['total']} — 누락 {', '.join(cov['missing'])}")
|
|
other = [i for i in issues if i["stage"] != "fetch"]
|
|
for issue in other:
|
|
lines.append(f"⚠️ {issue['stage']} 실패: {issue['error_type']}: {issue['message'][:200]}")
|
|
lines.append(f"{elapsed:.0f}초")
|
|
await activity_feed.post("✅ 수집 완료", lines)
|
|
|
|
|
|
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, "label": activity_feed.place_label(place)}
|
|
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"))
|
|
|
|
# 자동 해석 실패를 사장님에게 알리는 유일한 창구.
|
|
result["naver_place_missing"] = not any(
|
|
link.channel == LinkChannel.NAVER_PLACE.value for link in targets
|
|
)
|
|
|
|
await activity_feed.post(
|
|
"🔎 수집 시작",
|
|
[result["label"], f"place_id={place_id}",
|
|
f"크롤링 대상 {len(targets)}건: " + (", ".join(activity_feed.channel_name(l.channel) for l in targets) or "없음")],
|
|
)
|
|
|
|
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=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
|
|
channel_rows: list[str] = []
|
|
result["channels"] = channel_rows
|
|
|
|
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']}건 크롤링 생략")
|
|
channel_rows.append(f"필수 항목 확보 — 남은 {fetch_stat['skipped_enough']}건 생략")
|
|
break
|
|
|
|
name = activity_feed.channel_name(link.channel)
|
|
fetch_started = time.monotonic()
|
|
issues_before = len(collect_diagnostics.snapshot())
|
|
source, outcome = await fetch_one(link, PlaceCategory(place.category))
|
|
fetch_stat[outcome] += 1
|
|
if source is None:
|
|
fresh = collect_diagnostics.snapshot()[issues_before:]
|
|
reason = f"{fresh[-1]['error_type']}: {fresh[-1]['message'][:300]}" if fresh else outcome
|
|
channel_rows.append(f"❌ {name} — {reason}")
|
|
if outcome == "failed":
|
|
await activity_feed.post(
|
|
"⚠️ 채널 크롤링 실패",
|
|
[result["label"], f"{name} {link.url}", reason],
|
|
)
|
|
continue
|
|
|
|
# 수집 원문을 링크에 남긴다 — fact 가 아니라 '생성 근거' 자리다.
|
|
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"]
|
|
channel_rows.append(
|
|
f"✓ {name} — fact {f['stored']} · 사진 {m['stored']} · {time.monotonic() - fetch_started:.1f}초"
|
|
)
|
|
|
|
result["fetch"] = fetch_stat
|
|
result["units"] = unit_total
|
|
result["facts"] = facts_stat
|
|
result["media"] = media_stat
|
|
result["coverage"] = await coverage(place, place_id)
|
|
|
|
# 사진이 들어왔으면 분석을 이어서 건다 — 수집과 분석은 각각 몇 분이라 한 잡에 묶지 않는다.
|
|
if result["media"]["stored"] > 0:
|
|
result["vision_job_id"] = await _enqueue_vision(place_id, owner_user_id)
|
|
|
|
# 지역 데이터(주변 맛집·관광지·축제 + 지역 이야기)를 **여기서** 건다.
|
|
from services import story_service
|
|
result["local_job_id"] = await story_service.enqueue_region_job(place)
|
|
|
|
# ── 업소 조사 — 소개문을 쓸 재료 ──────────────────────────────────
|
|
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:
|
|
"""수집 중 채널이 알려준 예약 주소를 **예약 채널 링크**로 남긴다."""
|
|
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
|