import uuid from fastapi import Depends from common.category_schema import CategorySchemaError, get_schema from common.database.db_session_manager import DB_SESSION_MNG from common.database.model.models import place_facts, places, place_units from common.enums import ( FACT_STATUS_TRANSITIONS, LOCKED_FACT_STATUSES, PUBLISHABLE_FACT_STATUSES, DBWRType, ErrorType, FactStatus, FactWriteOutcome, PlaceCategory, SourceType, ) from common.logger import LOG from common.models.gmodel import UserInfo from common.utils.gtime import GTime from crud.fact_crud import FactCRUD, IFactCRUD from crud.place_crud import PlaceCRUD from router.v1.fact.protocol import ( FactData, FieldSpecData, Req_TransitionFact, Req_UpsertFact, Res_CategorySchema, Res_ExtractedFact, Res_ExtractFacts, Res_Fact, Res_FactList, ) from services.external.gemini_extract import extract_facts from services.intro_summary import summarize_intro from services.llm.gemini import GeminiError, GeminiNotConfigured # 사장님이 붙여넣은 원문의 출처 표기. OWNER_PASTE_SOURCE = "owner:paste" # 자동 출처. _AUTO_SOURCES = (SourceType.API, SourceType.CRAWL, SourceType.LLM) class FactService: """fact 기록 + 검증 상태 전이.""" def __init__(self, crud: IFactCRUD = Depends(FactCRUD), place_crud: PlaceCRUD = Depends(PlaceCRUD)): self.crud = crud self.place_crud = place_crud # 사업장 로드(회사 스코프) async def _load_place(self, user_info: UserInfo, place_id: str): err_type, place = await DB_SESSION_MNG.execute_lambda( places.DBType(), DBWRType.DB_READ.value, lambda s: self.place_crud.get_place(s, uuid.UUID(user_info.user_id), uuid.UUID(place_id)), ) if err_type != ErrorType.SUCCESS: return ErrorType.PLACE_NOT_FOUND, None return ErrorType.SUCCESS, place async def _mark_content_updated(self, place_id: str, ts): err = await DB_SESSION_MNG.execute_lambda_run( [places.DBType()], [lambda s: self._touch(s, place_id, ts)], ) if err != ErrorType.SUCCESS: LOG.e_no_callstack(f"[fact] content_updated_at 갱신 실패 place={place_id}") async def _touch(self, s, place_id: str, ts): from sqlalchemy import update query = update(places).where(places.place_id == uuid.UUID(place_id)).values(content_updated_at=ts, updated_at=ts) return await DB_SESSION_MNG.add(s, query) # 업종 스키마 노출(관리 화면 폼 생성용) async def get_category_schema(self, user_info: UserInfo, place_id: str) -> Res_CategorySchema: res = Res_CategorySchema() err_type, place = await self._load_place(user_info, place_id) if err_type != ErrorType.SUCCESS: res.result.SetResult(err_type) return res try: schema = get_schema(PlaceCategory(place.category)) except (CategorySchemaError, ValueError): res.result.SetResult(ErrorType.PLACE_INVALID_CATEGORY) return res res.category = schema.name res.label = schema.label res.fields = [FieldSpecData(**spec.to_dict()) for spec in schema.fields.values()] return res # 조회 async def list_facts(self, user_info: UserInfo, place_id: str, unit_id=None, publishable_only: bool = False) -> Res_FactList: res = Res_FactList() err_type, _place = await self._load_place(user_info, place_id) if err_type != ErrorType.SUCCESS: res.result.SetResult(err_type) return res list_err, rows = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: self.crud.list_facts(s, uuid.UUID(place_id), unit_id, None, publishable_only, True), ) if list_err != ErrorType.SUCCESS: res.result.SetResult(list_err) return res res.facts = [FactData.model_validate(r) for r in rows] await self._attach_summaries(res.facts) # 사이트에 나갈 수 있는 건수. res.publishable = sum(1 for r in rows if FactStatus(r.status) in PUBLISHABLE_FACT_STATUSES) # 재수집이 올려놓은 확인 대기 건수 — 관리 화면의 '검토할 것' 배지. res.pending_review = sum(1 for r in rows if FactStatus(r.status) == FactStatus.PENDING_OWNER) return res async def _attach_summaries(self, facts: list) -> None: """intro/room_intro 원문이 길면 캔버스 미리보기용 요약을 얹는다.""" for f in facts: f.summary = await summarize_intro(f.key, f.value) # 기록 async def extract_from_text(self, user_info: UserInfo, place_id: str, text: str) -> Res_ExtractFacts: """사장님이 붙여넣은 원문 → fact 후보.""" res = Res_ExtractFacts() err_type, place = await self._load_place(user_info, place_id) if err_type != ErrorType.SUCCESS: res.result.SetResult(err_type) return res try: schema = get_schema(PlaceCategory(place.category)) except (CategorySchemaError, ValueError): res.result.SetResult(ErrorType.PLACE_INVALID_CATEGORY) return res try: extracted = await extract_facts( place.name, PlaceCategory(place.category), text, source_url=OWNER_PASTE_SOURCE, ) except GeminiNotConfigured: res.result.SetResult(ErrorType.GENERATOR_NOT_CONFIGURED) return res except GeminiError as ex: LOG.w(f"[fact] 붙여넣기 추출 실패: {type(ex).__name__}: {ex}") res.result.SetResult(ErrorType.GENERATOR_CALL_FAILED) return res res.rejections = [[label, why] for label, why in extracted.rejected] res.rejected = len(extracted.rejected) # 단위(객실·메뉴) 이름을 실제 unit 으로 매핑한다. unit_map = await self._unit_map(place_id) for fact in extracted.facts: spec = schema.get(fact.key) row = Res_ExtractedFact( key=fact.key, label=(spec.label if spec else fact.key), value=str(fact.value or ""), unit_name=fact.unit_name, ) unit_id = unit_map.get(fact.unit_name) if fact.unit_name else None if fact.unit_name and unit_id is None: row.reason = f"'{fact.unit_name}' 단위가 아직 없다 — 수집으로 객실·메뉴가 먼저 만들어져야 한다" res.facts.append(row) res.rejected += 1 continue write = await self.upsert_fact(user_info, place_id, Req_UpsertFact( key=fact.key, value=fact.value, unit_id=unit_id, source_type=SourceType.LLM, source_url=OWNER_PASTE_SOURCE, )) if write.result.result == ErrorType.SUCCESS.value: row.stored = True res.stored += 1 else: row.reason = write.result.desc or "저장 실패" res.rejected += 1 res.facts.append(row) LOG.i(f"[fact] 붙여넣기 추출 — 저장 {res.stored}건 · 반려 {res.rejected}건 (place={place_id})") return res async def _unit_map(self, place_id: str) -> dict: """단위 이름 → unit_id.""" err, rows = await DB_SESSION_MNG.execute_lambda( place_units.DBType(), DBWRType.DB_READ.value, lambda s: self.place_crud.list_units(s, uuid.UUID(place_id)), ) if err != ErrorType.SUCCESS: return {} return {r.name: r.unit_id for r in (rows or []) if r.name} async def upsert_fact(self, user_info: UserInfo, place_id: str, req: Req_UpsertFact) -> Res_Fact: """fact 를 기록한다.""" res = Res_Fact() err_type, place = await self._load_place(user_info, place_id) if err_type != ErrorType.SUCCESS: res.result.SetResult(err_type) return res # 규칙 1 — 업종 스키마에 있는 key 만 try: schema = get_schema(PlaceCategory(place.category)) except (CategorySchemaError, ValueError): res.result.SetResult(ErrorType.PLACE_INVALID_CATEGORY) return res spec = schema.get(req.key) if spec is None: res.result.SetResult(ErrorType.FACT_INVALID_KEY) return res # 규칙 2 — owner 가 아니면 출처 URL 필수 if req.source_type != SourceType.OWNER and not (req.source_url or "").strip(): res.result.SetResult(ErrorType.FACT_SOURCE_REQUIRED) return res # 규칙 3 — ★ LLM 은 사실을 만들지 않는다. if req.source_type == SourceType.LLM and not spec.allow_llm: res.result.SetResult(ErrorType.FACT_INVALID_KEY) return res # 규칙 4 — TEMPLATE 은 FAQ 문의 안내 전용 출처다(services/faq_fill). if req.source_type == SourceType.TEMPLATE: res.result.SetResult(ErrorType.INVALID_REQUEST_DATA) return res pid = uuid.UUID(place_id) pub_err, published = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: self.crud.get_published_fact(s, pid, req.unit_id, req.key), ) if pub_err != ErrorType.SUCCESS: res.result.SetResult(pub_err) return res now = GTime.UTC() same_value = published is not None and (published.value or "") == (req.value or "") if req.source_type == SourceType.CRAWL and published is not None and ( published.source_type == SourceType.OWNER.value or FactStatus(published.status) in LOCKED_FACT_STATUSES ): if same_value: res.outcome = FactWriteOutcome.REFRESHED return await self._reload(res, pid, published.fact_id) return await self._write_candidate(res, pid, req, published, now, spec) # ── 값이 그대로다 — 검증을 초기화하지 않고 '언제 다시 확인했는지'만 갱신 ── if same_value: run_err, _rc = await DB_SESSION_MNG.execute_lambda_claim( place_facts.DBType(), lambda s: self.crud.refresh_collected(s, published.fact_id, req.source_type.value, req.source_url, now), ) if run_err != ErrorType.SUCCESS: res.result.SetResult(run_err) return res res.outcome = FactWriteOutcome.REFRESHED return await self._reload(res, pid, published.fact_id) # LLM 이 쓴 문장은 후보로 두지 않고 **바로 노출값**이다. if req.source_type in (SourceType.LLM, SourceType.CRAWL): if published is not None and FactStatus(published.status) in LOCKED_FACT_STATUSES: return await self._write_candidate(res, pid, req, published, now, spec) return await self._replace_published(res, place_id, pid, req, published, now, spec, user_info) if req.source_type in _AUTO_SOURCES: return await self._write_candidate(res, pid, req, published, now, spec) return await self._replace_published(res, place_id, pid, req, published, now, spec, user_info) async def _write_candidate(self, res, pid, req, published, now, spec): """자동 수집 — 노출값은 건드리지 않고 후보로 적재한다.""" target_status = FactStatus.PENDING_OWNER if published is not None else FactStatus.UNVERIFIED cand_err, candidate = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: self.crud.get_candidate(s, pid, req.unit_id, req.key, req.source_type.value), ) if cand_err != ErrorType.SUCCESS: res.result.SetResult(cand_err) return res # 같은 출처가 이미 올려둔 후보가 있으면 갱신한다(같은 후보가 계속 쌓이지 않게). if candidate is not None: run_err, _rc = await DB_SESSION_MNG.execute_lambda_claim( place_facts.DBType(), lambda s: self.crud.update_candidate( s, candidate.fact_id, req.value, req.source_url, target_status.value, now ), ) if run_err != ErrorType.SUCCESS: res.result.SetResult(run_err) return res res.outcome = FactWriteOutcome.CANDIDATE_UPDATED return await self._reload(res, pid, candidate.fact_id) fact = place_facts( place_id=pid, unit_id=req.unit_id, key=req.key, value=req.value, unit=spec.unit, source_type=req.source_type.value, source_url=(req.source_url or None), collected_at=now, status=target_status.value, expires_at=req.expires_at, ) run_err = await DB_SESSION_MNG.execute_lambda_run( [place_facts.DBType()], [lambda s: self.crud.add_fact(s, fact)], ) if run_err != ErrorType.SUCCESS: res.result.SetResult(run_err) return res res.outcome = FactWriteOutcome.CANDIDATE_CREATED res.fact = FactData.model_validate(fact) return res async def _replace_published(self, res, place_id, pid, req, published, now, spec, user_info): """직접 입력·크롤링·허용된 LLM 문장을 노출값으로 교체한다.""" fact = place_facts( place_id=pid, unit_id=req.unit_id, key=req.key, value=req.value, unit=spec.unit, source_type=req.source_type.value, source_url=(req.source_url or None), collected_at=now, verified_by=None if req.source_type == SourceType.CRAWL else uuid.UUID(user_info.user_id), verified_at=now, status=FactStatus.VERIFIED.value, expires_at=req.expires_at, ) run_err = await DB_SESSION_MNG.execute_lambda_run( [place_facts.DBType()], [ lambda s: self._expire_then_ok(s, pid, req, now), lambda s: self.crud.add_fact(s, fact), *([lambda s: self._reject_others_ok(s, pid, fact, now, fact.fact_id)] if req.source_type == SourceType.CRAWL else []), ], ) if run_err != ErrorType.SUCCESS: res.result.SetResult(run_err) return res await self._mark_content_updated(place_id, now) res.outcome = FactWriteOutcome.PUBLISHED_REPLACED if published is not None else FactWriteOutcome.PUBLISHED_CREATED res.fact = FactData.model_validate(fact) return res async def _expire_then_ok(self, s, pid, req, now): """execute_lambda_run 은 각 람다가 ErrorType 을 돌려주길 요구한다 — rowcount 는 버린다.""" err_type, _rowcount = await self.crud.expire_published(s, pid, req.unit_id, req.key, now) return err_type async def _reload(self, res, pid, fact_id): _e, row = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: self.crud.get_fact(s, pid, fact_id), ) res.fact = FactData.model_validate(row) if row is not None else None return res # 검증 상태 전이 async def transition(self, user_info: UserInfo, place_id: str, fact_id: str, req: Req_TransitionFact) -> Res_Fact: """검증 상태를 바꾼다.""" res = Res_Fact() err_type, _place = await self._load_place(user_info, place_id) if err_type != ErrorType.SUCCESS: res.result.SetResult(err_type) return res pid = uuid.UUID(place_id) fid = uuid.UUID(fact_id) get_err, fact = await DB_SESSION_MNG.execute_lambda( place_facts.DBType(), DBWRType.DB_READ.value, lambda s: self.crud.get_fact(s, pid, fid), ) if get_err != ErrorType.SUCCESS: res.result.SetResult(ErrorType.FACT_NOT_FOUND) return res current = FactStatus(fact.status) target = req.status if target not in FACT_STATUS_TRANSITIONS.get(current, set()): res.result.SetResult(ErrorType.FACT_INVALID_TRANSITION) return res now = GTime.UTC() data: dict = {} if target in PUBLISHABLE_FACT_STATUSES: data["verified_by"] = uuid.UUID(user_info.user_id) data["verified_at"] = now if target == FactStatus.CORRECTED: if req.value is None: res.result.SetResult(ErrorType.INVALID_REQUEST_DATA) return res # 사람이 고친 값 — 출처가 owner 로 바뀌고 자동 수집은 이후 후보로만 도전할 수 있다. data["value"] = req.value data["source_type"] = SourceType.OWNER.value promoting = target in PUBLISHABLE_FACT_STATUSES and current not in PUBLISHABLE_FACT_STATUSES funcs = [] if promoting: # 승격 전에 자리를 비운다 — 노출값 유니크(1건) 때문에 순서가 중요하다. funcs.append(lambda s: self._expire_published_ok(s, pid, fact, now)) funcs.append(lambda s: self._transition_ok(s, fid, current, target, data)) if promoting: funcs.append(lambda s: self._reject_others_ok(s, pid, fact, now, fid)) run_err = await DB_SESSION_MNG.execute_lambda_run([place_facts.DBType()], funcs) if run_err != ErrorType.SUCCESS: res.result.SetResult( ErrorType.FACT_INVALID_TRANSITION if run_err == ErrorType.DB_EMPTY_DATA else run_err ) return res # 노출값이 바뀌는 전이(승격 / 정정 / 노출값 내림)면 재빌드 대상으로 표시한다. if promoting or current in PUBLISHABLE_FACT_STATUSES: await self._mark_content_updated(place_id, now) res.outcome = FactWriteOutcome.PUBLISHED_REPLACED if promoting else None return await self._reload(res, pid, fid) async def _expire_published_ok(self, s, pid, fact, now): err_type, _rc = await self.crud.expire_published(s, pid, fact.unit_id, fact.key, now, except_fact_id=fact.fact_id) return err_type async def _reject_others_ok(self, s, pid, fact, now, keep_fact_id): err_type, _rc = await self.crud.reject_candidates(s, pid, fact.unit_id, fact.key, now, except_fact_id=keep_fact_id) return err_type async def _transition_ok(self, s, fid, current, target, data): err_type, rowcount = await self.crud.transition(s, fid, (current.value,), target.value, data) if err_type != ErrorType.SUCCESS: return err_type return ErrorType.SUCCESS if rowcount else ErrorType.DB_EMPTY_DATA