"""기존 크롤링 후보를 현재 기록 정책으로 다시 적용한다. --apply 없이는 조회만 한다.""" import argparse import asyncio import json import sys import uuid from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[1])) from sqlalchemy import select from common.database.db_session_manager import DB_SESSION_MNG from common.database.model.models import place_facts, places from common.enums import DBWRType, ErrorType, FactStatus, SourceType from common.models.gmodel import UserInfo from crud.fact_crud import FactCRUD from crud.place_crud import PlaceCRUD from router.v1.fact.protocol import Req_UpsertFact from services.fact_service import FactService async def main(place_id, apply): try: err, rows = await DB_SESSION_MNG.execute_lambda( places.DBType(), DBWRType.DB_READ.value, lambda s: DB_SESSION_MNG.execute(s, select(places).where( places.place_id == place_id, places.deleted.is_(False))), ) if err != ErrorType.SUCCESS or len(rows or []) != 1: raise RuntimeError('사업장 조회 실패') place = rows[0] err, facts = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: DB_SESSION_MNG.execute(s, select(place_facts).where( place_facts.place_id == place_id, place_facts.deleted.is_(False), place_facts.source_type == SourceType.CRAWL.value, place_facts.status.in_([FactStatus.UNVERIFIED.value, FactStatus.PENDING_OWNER.value]), ).order_by(place_facts.collected_at.desc())), ) if err != ErrorType.SUCCESS: raise RuntimeError('크롤링 후보 조회 실패') service = FactService(FactCRUD(), PlaceCRUD()) actor = UserInfo(user_id=str(place.owner_user_id), id='crawl-policy', role=1) seen = set() result = [] for fact in facts or []: key = (fact.unit_id, fact.key) if key in seen: continue seen.add(key) row = {'unitId': str(fact.unit_id), 'key': fact.key, 'value': fact.value} if apply: res = await service.upsert_fact(actor, str(place_id), Req_UpsertFact( key=fact.key, value=fact.value, unit_id=fact.unit_id, source_type=SourceType.CRAWL, source_url=fact.source_url, expires_at=fact.expires_at, )) if not res.result.success: raise RuntimeError(f'{fact.key}: {res.result.desc}') row['status'] = res.fact.status result.append(row) print(json.dumps({'applied': apply, 'place': place.name, 'facts': result}, ensure_ascii=False, default=str)) finally: await DB_SESSION_MNG.dispose_all() if __name__ == '__main__': parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('--place-id', required=True, type=uuid.UUID) parser.add_argument('--apply', action='store_true') args = parser.parse_args() asyncio.run(main(args.place_id, args.apply))