o2o-site-AEO/solution/backend/services/collect_service.py
Mina Choi 9d25ed613e 구조: 사장님(solution)과 내부 운영(admin)을 두 앱으로 가른다
최상단을 프로젝트 단위로 평평하게 둔다 — o2o-negosium 과 같은 규약이고, 이 레포만
다르게 갈 이유가 없다. negodata/{backend,front} 가 프로젝트 안에서 f/b 를 가르는 선례,
lps-admin/ 이 백엔드 없이 프론트만 가진 최상단 폴더의 선례다.

  backend/ frontend/{admin,site,shared}  →  solution/{backend,front,site,shared} + admin/

## 왜

내부 라우트(/local-content, /places/:id/seo)의 이름과 화면 코드가 사장님 번들에
그대로 실려 나가고 있었다. UserRole.DEVELOPER 주석의 "고객사에 존재를 노출하지 않는다"를
번들이 깨고 있었다 — 라우트 가드는 화면을 가리지 번들은 못 가린다.
번들을 갈라 확인했다: 사장님 dist 에서 local-content · /places · SeoAudit 이 전부 0건이다.

그 과정에서 두 곳이 더 새고 있었다.
- AppShell 의 NAV 배열이 내부 메뉴를 하드코딩하고 있었다. 앱을 가른 뒤에도 dist 에
  local-content 가 남아서 찾았다. 메뉴는 이제 앱이 prop 으로 들고 온다.
- EditorHeader·BuilderPage·LoginPage 가 /places 로 링크하고 있었다. 그 화면이 admin 으로
  나갔으니 사장님 앱에서는 404 다. 링크를 걷어내고 LoginPage 기본 도착지는 '/' 로 바꿨다
  (앱마다 홈이 다르고 각 라우터의 '/' 가 이미 그걸 안다).

## admin 에 백엔드를 두지 않았다

내부 화면이 부르는 훅이 전부 router/v1/{place,fact,local,validator} 에 이미 있다.
자체 백엔드를 두면 place·fact·link 를 같은 DB 에 대고 두 번 구현하게 된다.
대가는 solution/backend 가 죽으면 admin 도 멈추는 것 — 내부 도구라 감수한다.

## admin 의 `@` 는 solution/front/src 를 가리킨다

내부 화면이 쓰는 API 클라이언트·UI·수집 배선이 solution 에 한 벌만 있고 그 파일들끼리도
`@/...` 로 서로를 부른다. admin 에서 `@` 를 자기 src 로 잡으면 그 참조가 전부 깨진다
(실측 TS2307 14건). 복제하는 길도 있지만 RecollectPanel 주석이 금지한다 —
"수집 경로를 두 벌 만들면 확정 게이트"가 갈라진다.
admin 자기 파일만 `@admin` 이고, 의존 방향은 admin → solution 한 쪽뿐이다.

admin 이 여는 빌더는 다른 오리진이라 절대 URL + 새 탭이다(admin/src/lib/solutionUrl.ts).
react-router Link 로 두면 admin 안에서 라우트를 찾다 404 다.

## 그 밖

- npm 워크스페이스 루트를 레포 루트로 올렸다(admin 이 solution 밖이라).
- docker-compose 를 255→174줄로 줄이고 admin(:3002) 서비스를 넣었다. ADMIN_BIND 기본값은
  127.0.0.1 — 0.0.0.0 으로 열면 앱을 가른 의미가 없다.
- 발행 호스트를 프론트 .env 에 따로 적지 않는다. compose 가 루트의 SITE_PUBLIC_HOST 를
  VITE_PUBLISH_HOST 로 흘려보낸다 — 두 곳에 적으면 canonical 과 화면 주소가 조용히 갈라진다.
- nginx/site.conf 를 git 에서 빼고 .example 만 남겼다(.env·*.toml 과 같은 규약).
  compose 가 bind mount 하므로 클론 직후 복사해야 한다 — 없으면 Docker 가 그 자리에
  디렉토리를 만들어 nginx 가 설정 없이 뜬다.
- config.test.toml.example 을 추가했다. 없으면 클론한 사람이 pytest 를 아예 못 돌린다
  (conftest import 단계에서 죽는다). 외부 API 키는 전부 빈값이다 —
  APP_ENV=test 가 .env 를 안 읽는 이유를 여기서 우회하면 안 된다.
- 경로가 한 칸 깊어져 test_schema_ddl(parents[2]→[3]) 과 test_site_theme 을 고쳤다.

검증: front·admin·site 전부 lint 0 / build 0. 백엔드 514 passed.
남은 4건(test_build_publish 3 · test_snapshot 1)은 이 변경 전부터 실패하던 것으로,
손대지 않은 메인 체크아웃에서 같은 4건이 같게 실패하는 것을 확인했다.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019uYhHQdssRubirPirrdJJC
2026-08-31 15:12:09 +09:00

574 lines
28 KiB
Python

"""채널 발견 → 확정 URL 크롤링 → fact·사진 후보 저장.
수집 결과는 사용자가 승인하기 전까지 사이트에 노출하지 않는다.
"""
import uuid
from common.database.db_session_manager import DB_SESSION_MNG
from common.database.model.models import facts as facts_model
from common.database.model.models import media, place_links, places, 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.external import naver_place_lookup, perplexity, tour_lookup
from services.fact_service import FactService
from router.v1.fact.protocol import Req_UpsertFact
_place_crud = PlaceCRUD()
_fact_crud = FactCRUD()
class CollectAborted(RuntimeError):
"""재시도해도 소용없는 중단 — 잡의 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_links.DBType()],
[lambda s: _place_crud.add_link(s, place_links(
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_links.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_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 — 발견 실패가 수집을 죽이면 안 된다
LOG.w(f"[collect] TourAPI 조회 실패(계속): {type(ex).__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_links.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"
async def discover_links(place, place_id: str, *, include_perplexity: bool = False) -> dict:
"""채널 URL 을 찾아 place_links 에 적재한다(미확정 상태).
★ 순서가 곧 신뢰도다. 지금 자동 발견 경로는 **네이버 플레이스 직접 해석 하나뿐**이고,
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 — 발견 실패가 수집을 죽이면 안 된다
LOG.w(f"[collect] 네이버 플레이스 조회 실패(계속): {type(ex).__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
LOG.w(f"[collect] TourAPI 조회 실패(계속): {type(ex).__name__}: {ex}")
stat["tour_api"] = "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:
# ★ 실패해도 파이프라인을 죽이지 않는다 — 이미 등록된 링크로 크롤링은 계속한다.
LOG.w(f"[collect] URL 발견 실패(계속): {type(ex).__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_links(
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_links.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_links.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_links.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):
"""링크 하나를 긁는다. 실패해도 예외를 던지지 않는다 — 나머지 링크가 살아야 한다."""
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)
except Exception as ex:
LOG.w(f"[collect] 수집 실패(계속) {link.url}: {type(ex).__name__}: {ex}")
return None, "failed"
if not source.ok:
LOG.w(f"[collect] 수집 실패(계속) {link.url}: {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(
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 = units(place_id=uuid.UUID(place_id), name=name, sort_order=order)
run_err = await DB_SESSION_MNG.execute_lambda_run(
[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(
media.DBType(),
DBWRType.DB_READ.value,
lambda s: DB_SESSION_MNG.execute(
s, select(media).where(media.place_id == uuid.UUID(place_id), media.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 = media(
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(
[media.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 에 저장돼 폴링·감사에 쓰인다."""
payload = job["payload"]
place_id = payload["place_id"]
company_id = payload["company_id"]
err, place = await DB_SESSION_MNG.execute_lambda(
places.DBType(),
DBWRType.DB_READ.value,
lambda s: _place_crud.get_place(s, uuid.UUID(company_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, company_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, company_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=payload.get("requested_by") or str(place.verified_by or uuid.uuid4()),
id="collector",
company_id=company_id,
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)
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_links.DBType(),
lambda sess, u=link.url, t=source.text: _place_crud.set_link_raw(
sess, uuid.UUID(place_id), u, {"text": t[:8000]},
),
)
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, company_id)
await _finish(place_id, company_id, PlaceStatus.REVIEW)
LOG.i(f"[collect] 완료 place={place_id} fact {result['facts']['stored']}건 · 사진 {result['media']['stored']}장")
return result
async def _finish(place_id: str, company_id: str, status: PlaceStatus):
"""수집이 끝나면 사업장을 검수 대기로 돌린다 — 수집값은 전부 후보라 사람이 봐야 한다."""
await DB_SESSION_MNG.execute_lambda_claim(
places.DBType(),
lambda s: _place_crud.update_place(
s, uuid.UUID(company_id), uuid.UUID(place_id), {"status": status.value}
),
)
async def _enqueue_vision(place_id: str, company_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("[collect] GEMINI_API_KEY 미설정 — 사진 분석 건너뜀(사진은 확인 큐에 남는다)")
return None
job_id, _created = await enqueue_job(
JobQueue(), JobType.VISION,
{"place_id": place_id, "company_id": company_id},
dedupe_key=f"vision:{place_id}",
)
return job_id