# backend/app/connectors/base.py # 도메인별 커넥터 추상 인터페이스. mock=시드 기반, real=외부 제공자(phase-13). import hashlib from abc import ABC, abstractmethod from collections.abc import Iterable from dataclasses import dataclass, field from datetime import UTC, datetime from typing import Any, Optional from sqlmodel import Session, select from ..models import ConnectorAccount, ConnectorSyncState, ConnState, ExternalLink class CalendarConnector(ABC): """calendar 도메인 읽기-전용 커넥터(phase-8). mock=시드(cal-data.js), real=Google Calendar. 페이지 읽기는 이 인터페이스만 의존. (phase-13 의 sync 프레임워크와 공존).""" @abstractmethod def list_events(self, session, day: Optional[int] = None) -> list[dict]: ... @abstractmethod def get_meeting(self, session, meeting_id: str) -> Optional[dict]: ... @abstractmethod def write_event(self, session, payload: dict) -> dict: ... # ============================================================================ # phase-13: 통합 sync 프레임워크 — fetch / normalize / write / sync (인터페이스 불변) # 위 읽기-전용 커넥터와 공존한다. 이쪽은 "수집(ingestion)" 계약. # ============================================================================ def _now() -> datetime: return datetime.now(UTC) @dataclass class RawRecord: """제공자에서 fetch 한 1건의 원시 레코드(provider 페이로드).""" external_id: str payload: dict etag: str = "" external_updated_at: Optional[datetime] = None @dataclass class NormalizedRecord: """normalize() 산출물 — 내부 모델로 write 가능한 정규화 dict.""" entity_type: str # "email" | "event" | "finance_tx" | ... external_id: str fields: dict[str, Any] = field(default_factory=dict) etag: str = "" external_updated_at: Optional[datetime] = None @dataclass class SyncResult: domain: str account_id: str seen: int = 0 upserted: int = 0 skipped: int = 0 errors: int = 0 events_published: list[str] = field(default_factory=list) detail: str = "" class BaseConnector(ABC): """모든 도메인 커넥터의 공통 계약. Mock/Real 이 동일하게 구현한다(인터페이스 불변).""" domain: str = "base" entity_type: str = "record" def __init__(self, account: ConnectorAccount): self.account = account # ── (1) fetch: 제공자에서 원시 레코드를 가져온다(증분 since/cursor) ── @abstractmethod def fetch(self, session: Session, *, full: bool = False) -> Iterable[RawRecord]: ... # ── (2) normalize: provider 페이로드 → 내부 모델 dict ── @abstractmethod def normalize(self, raw: RawRecord) -> NormalizedRecord: ... # ── (3) write: 정규화 레코드를 내부 테이블에 upsert(external_link 멱등) ── @abstractmethod def write(self, session: Session, norm: NormalizedRecord) -> tuple[str, bool]: """return (entity_id, created) — created=False 면 update/skip.""" ... # ── 도메인이 발행할 연합 이벤트 타입(없으면 None) ── def event_for(self, norm: NormalizedRecord, entity_id: str) -> Optional[str]: return None # ── (4) sync: fetch→normalize→write 오케스트레이션 + event_bus 발행 ── def sync(self, session: Session, *, full: bool = False) -> SyncResult: from ..automation.event_bus import bus # phase-7 정본 (지연 import) res = SyncResult(domain=self.domain, account_id=self.account.id) self._set_state(session, ConnState.syncing) try: for raw in self.fetch(session, full=full): res.seen += 1 if self._is_duplicate(session, raw): # etag 동일 → skip(멱등) res.skipped += 1 continue norm = self.normalize(raw) # provider 가 스팸으로 분류한 메일은 저장하지 않는다(유저: 아예 안 가져옴). # folder 없는 도메인(캘린더 등)은 .get→None 이라 영향 없음. if norm.fields.get("folder") == "spam": res.skipped += 1 continue entity_id, created = self.write(session, norm) self._upsert_link(session, raw, norm, entity_id) if created: res.upserted += 1 evt = self.event_for(norm, entity_id) if evt: bus.publish( evt, { "account_id": self.account.id, "entity_type": norm.entity_type, "entity_id": entity_id, }, ) res.events_published.append(evt) else: res.skipped += 1 session.commit() self._finish_ok(session, res) bus.publish( "connector.synced", { "account_id": self.account.id, "upserted": res.upserted, "skipped": res.skipped, "errors": res.errors, }, ) except Exception as e: # 부분 실패도 서버 죽지 않음(오프라인 폴백 철학) session.rollback() res.errors += 1 res.detail = f"{type(e).__name__}: {e}" self._finish_error(session, res.detail) bus.publish( "connector.error", {"account_id": self.account.id, "detail": res.detail, "state": "error"}, ) return res # ── 공통: 중복/멱등 ── def _link_for(self, session: Session, external_id: str) -> Optional[ExternalLink]: return session.exec( select(ExternalLink).where( ExternalLink.account_id == self.account.id, ExternalLink.external_id == external_id, ) ).first() def _is_duplicate(self, session: Session, raw: RawRecord) -> bool: link = self._link_for(session, raw.external_id) return bool(link and raw.etag and link.etag == raw.etag) def _upsert_link(self, session, raw, norm, entity_id): link = self._link_for(session, raw.external_id) if link: link.etag = norm.etag link.entity_id = entity_id link.external_updated_at = norm.external_updated_at session.add(link) else: xlid = ( "xl-" + hashlib.sha1(f"{self.account.id}:{raw.external_id}".encode()).hexdigest()[:12] ) session.add( ExternalLink( id=xlid, account_id=self.account.id, external_id=raw.external_id, entity_type=norm.entity_type, entity_id=entity_id, etag=norm.etag, external_updated_at=norm.external_updated_at, ) ) # ── 상태/커서/로그 ── def _state_row(self, session) -> ConnectorSyncState: st = session.get(ConnectorSyncState, f"cs-{self.account.id}") if not st: st = ConnectorSyncState(id=f"cs-{self.account.id}", account_id=self.account.id) session.add(st) return st def _set_state(self, session, state: ConnState): self.account.state = state self.account.updated_at = _now() session.add(self.account) session.commit() def _finish_ok(self, session, res: SyncResult): self.account.state = ConnState.connected self.account.last_synced_at = _now() if res.upserted or res.seen: self.account.last_label = "방금 동기화" self.account.error_detail = "" st = self._state_row(session) st.items_seen += res.seen st.items_upserted += res.upserted st.last_delta_sync_at = _now() session.add_all([self.account, st]) session.commit() self._log(session, "sync", f"upserted={res.upserted} skipped={res.skipped}", res.upserted) def _finish_error(self, session, detail: str): self.account.state = ConnState.error self.account.error_detail = detail self.account.last_label = "동기화 실패 · 다시 시도" session.add(self.account) session.commit() self._log(session, "error", detail) def _log(self, session, action: str, detail: str = "", items: int = 0): import uuid from ..models import ConnectorAccountLog session.add( ConnectorAccountLog( id="cl-" + uuid.uuid4().hex[:8], account_id=self.account.id, action=action, detail=detail, items=items, ) ) session.commit()