Зачем коннекторы и чем они отличаются от загрузчиков

Загрузчик из предыдущего урока работает с локальным файлом: открыл, прочитал, закрыл. Коннектор решает другую задачу — он работает с живым удалённым источником:

  • Аутентификация — нужен токен, который может протухнуть.
  • Пагинация — API отдаёт по 100 записей, а страниц тысячи.
  • Rate limiting — нельзя слать 1000 запросов в секунду.
  • Иерархия — страница вложена в раздел, раздел — в пространство.
  • Инкрементальная синхронизация — при повторном запуске брать только изменения.
  • Разные форматы ответа — Notion отдаёт блоки JSON, Confluence — HTML, Drive — файлы.

На выходе коннектор, как и загрузчик, должен отдавать list[Document] — единый интерфейс для остального пайплайна.

100%
колёсико — масштаб  ·  зажать и тянуть — перемещение
ИСТОЧНИКИ Notion Block API · pages · DBs Confluence REST v2 · spaces · pages Google Drive Drive API · Docs API SharePoint Graph API · drives · lists BaseConnector Auth / Token Pagination Rate Limiter Incremental sync → list[Document] Document page_content · metadata source · url · last_edited · space · … → Chunking → Embedding → Vector DB Sync State (cursor) last_sync_at · processed_ids · etag API token / OAuth 2.0 Rate: 3–10 req/s Cursor pagination Retry on 429/503

Аутентификация: токены, OAuth и сервисные аккаунты

Первое, что нужно сделать при работе с любым API — разобраться с аутентификацией. Каждый сервис использует свой метод, и ошибки здесь приводят к самым непонятным падениям.

Сервис
Метод
Как получить
Rate limit
Notion
API key (Internal)
Settings → Integrations → «+ New integration». Токен начинается с secret_. Страницы нужно явно расшарить интеграции.
3 req/s
Confluence
Basic Auth / PAT
Atlassian API Token (Account Settings → Security). Передаётся как Basic Auth: email:token в base64. Cloud и Server — разные URL.
10 req/s
Google Drive
Service Account
Google Cloud Console → IAM → Service Account → JSON key. Нужно расшарить папки на email сервисного аккаунта. OAuth для пользовательских данных.
1000 req/100s
SharePoint
OAuth 2.0 (MSAL)
Azure AD → App Registration → Client Secret. Нужны права Sites.Read.All или Files.Read.All. Client Credentials flow для серверных приложений.
10 000 req/10min

OAuth 2.0 Client Credentials: как это работает

SharePoint (и многие корпоративные API) используют OAuth 2.0 в режиме Client Credentials — это flow для серверных приложений без участия пользователя. Понимать его важно, чтобы правильно настроить Azure AD и не путаться с токенами.

1
Регистрация приложения в Azure AD. Получаем client_id и tenant_id. Создаём client_secret (действует 1–2 года, нужно обновлять).
2
Запрос access token: POST на https://login.microsoftonline.com/{tenant_id}/oauth2/v2.0/token с параметрами grant_type=client_credentials, scope=https://graph.microsoft.com/.default. В ответ — JWT-токен со сроком жизни 3600 секунд.
3
Запросы к API с заголовком Authorization: Bearer <access_token>. Токен кэшируем — не запрашиваем новый на каждый вызов.
4
Обновление токена: проверяем время жизни перед запросом. MSAL делает это автоматически при использовании ConfidentialClientApplication.

Базовый класс коннектора

Прежде чем писать под каждый сервис, вынесем общую логику: rate limiting, retry с exponential backoff и пагинацию. Эти паттерны одинаковы для всех API.

# connectors/base.py
import asyncio
import time
import logging
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import AsyncIterator
from pathlib import Path

import httpx
from core.document import Document

logger = logging.getLogger(__name__)


@dataclass
class SyncState:
    """
    Состояние синхронизации — сохраняем между запусками,
    чтобы следующий раз брать только изменения.
    """
    last_sync_at: str = ""          # ISO 8601 timestamp
    processed_ids: set[str] = field(default_factory=set)
    cursor: str = ""                # cursor/page_token для resume

    def save(self, path: Path) -> None:
        import json
        data = {
            "last_sync_at": self.last_sync_at,
            "processed_ids": list(self.processed_ids),
            "cursor": self.cursor,
        }
        path.write_text(json.dumps(data, ensure_ascii=False, indent=2))

    @classmethod
    def load(cls, path: Path) -> "SyncState":
        import json
        if not path.exists():
            return cls()
        data = json.loads(path.read_text())
        return cls(
            last_sync_at=data.get("last_sync_at", ""),
            processed_ids=set(data.get("processed_ids", [])),
            cursor=data.get("cursor", ""),
        )


class RateLimiter:
    """Токен-баккет: пропускает не более max_rps запросов в секунду."""

    def __init__(self, max_rps: float = 3.0):
        self._interval = 1.0 / max_rps
        self._last_call = 0.0

    async def acquire(self) -> None:
        now = time.monotonic()
        wait = self._interval - (now - self._last_call)
        if wait > 0:
            await asyncio.sleep(wait)
        self._last_call = time.monotonic()


class BaseConnector(ABC):
    """
    Базовый класс для всех коннекторов.
    Предоставляет: rate limiting, retry, общий HTTP-клиент.
    """

    MAX_RETRIES = 3
    RETRY_STATUSES = {429, 500, 502, 503, 504}

    def __init__(self, max_rps: float = 3.0):
        self._limiter = RateLimiter(max_rps)
        self._client: httpx.AsyncClient | None = None

    async def __aenter__(self):
        self._client = httpx.AsyncClient(timeout=30)
        return self

    async def __aexit__(self, *args):
        if self._client:
            await self._client.aclose()

    async def _get(self, url: str, **kwargs) -> dict:
        """GET с rate limiting и retry."""
        await self._limiter.acquire()
        for attempt in range(self.MAX_RETRIES):
            response = await self._client.get(url, **kwargs)
            if response.status_code == 200:
                return response.json()
            if response.status_code in self.RETRY_STATUSES:
                # 429 может содержать Retry-After заголовок
                retry_after = float(
                    response.headers.get("Retry-After", 2 ** attempt)
                )
                logger.warning(
                    f"HTTP {response.status_code} from {url}, "
                    f"retry {attempt + 1}/{self.MAX_RETRIES} "
                    f"after {retry_after:.1f}s"
                )
                await asyncio.sleep(retry_after)
            else:
                response.raise_for_status()
        raise RuntimeError(f"Max retries exceeded for {url}")

    async def _post(self, url: str, **kwargs) -> dict:
        """POST с rate limiting и retry (нужен для Notion API)."""
        await self._limiter.acquire()
        for attempt in range(self.MAX_RETRIES):
            response = await self._client.post(url, **kwargs)
            if response.status_code in (200, 201):
                return response.json()
            if response.status_code in self.RETRY_STATUSES:
                await asyncio.sleep(2 ** attempt)
            else:
                response.raise_for_status()
        raise RuntimeError(f"Max retries exceeded for {url}")

    @abstractmethod
    async def load(self, state: SyncState | None = None) -> list[Document]:
        """Загрузить документы. state=None → полная загрузка."""
        ...

Notion: блочная модель контента

Notion хранит содержимое страниц не как текст, а как дерево блоков. Страница — корневой блок. Внутри — параграфы, заголовки, списки, таблицы, вложенные страницы — каждый элемент является отдельным блоком с типом. Это мощная модель для редактирования, но для RAG её нужно «развернуть» в плоский текст.

Page (block type: "page")
├── heading_1:  "Введение"
├── paragraph:  "Текст вводного раздела…"
├── bulleted_list_item: "Пункт 1"
├── bulleted_list_item: "Пункт 2"
├── child_page: "Вложенная страница"  ← рекурсивно
└── table
    ├── table_row: ["Имя", "Значение"]
    └── table_row: ["foo",  "bar"]
# connectors/notion_connector.py
import os
from datetime import datetime, timezone
from core.document import Document
from connectors.base import BaseConnector, SyncState


NOTION_VERSION = "2022-06-28"

# Типы блоков, которые содержат текст
TEXT_BLOCK_TYPES = {
    "paragraph", "heading_1", "heading_2", "heading_3",
    "bulleted_list_item", "numbered_list_item", "toggle",
    "quote", "callout", "code",
}

HEADING_PREFIX = {
    "heading_1": "# ", "heading_2": "## ", "heading_3": "### ",
}


def _extract_rich_text(rich_text_list: list) -> str:
    """Извлекаем текст из Notion rich_text массива."""
    return "".join(rt.get("plain_text", "") for rt in rich_text_list)


def _block_to_text(block: dict) -> str:
    """Преобразуем блок в строку Markdown."""
    btype = block.get("type", "")
    content = block.get(btype, {})

    if btype in TEXT_BLOCK_TYPES:
        text = _extract_rich_text(content.get("rich_text", []))
        prefix = HEADING_PREFIX.get(btype, "")
        if btype == "bulleted_list_item":
            prefix = "- "
        elif btype == "code":
            lang = content.get("language", "")
            return f"```{lang}\n{text}\n```"
        return f"{prefix}{text}" if text else ""

    if btype == "table_row":
        cells = [
            _extract_rich_text(cell)
            for cell in content.get("cells", [])
        ]
        return "| " + " | ".join(cells) + " |"

    return ""  # image, divider, embed — пропускаем


class NotionConnector(BaseConnector):
    """
    Загружает страницы Notion.
    Обходит все дочерние блоки рекурсивно, собирает текст,
    поддерживает инкрементальную синхронизацию по last_edited_time.
    """

    BASE_URL = "https://api.notion.com/v1"

    def __init__(self, token: str | None = None, max_rps: float = 3.0):
        super().__init__(max_rps=max_rps)
        self._token = token or os.environ["NOTION_TOKEN"]

    @property
    def _headers(self) -> dict:
        return {
            "Authorization": f"Bearer {self._token}",
            "Notion-Version": NOTION_VERSION,
            "Content-Type": "application/json",
        }

    async def load(self, state: SyncState | None = None) -> list[Document]:
        """Загружает все расшаренные с интеграцией страницы."""
        pages = await self._search_all_pages(
            edited_after=state.last_sync_at if state else ""
        )
        documents = []
        for page in pages:
            page_id = page["id"]
            if state and page_id in state.processed_ids:
                continue
            doc = await self._load_page(page)
            if doc:
                documents.append(doc)
        return documents

    async def load_page(self, page_id: str) -> Document | None:
        """Загружает одну страницу по ID."""
        page = await self._get(
            f"{self.BASE_URL}/pages/{page_id}",
            headers=self._headers,
        )
        return await self._load_page(page)

    async def _search_all_pages(self, edited_after: str = "") -> list[dict]:
        """Получаем все страницы через Search API с пагинацией."""
        pages = []
        cursor = None
        payload: dict = {
            "filter": {"property": "object", "value": "page"},
            "page_size": 100,
            "sort": {"direction": "descending", "timestamp": "last_edited_time"},
        }
        if edited_after:
            # Notion Search API не поддерживает фильтр по дате напрямую,
            # фильтруем на стороне клиента
            pass

        while True:
            if cursor:
                payload["start_cursor"] = cursor
            data = await self._post(
                f"{self.BASE_URL}/search",
                headers=self._headers,
                json=payload,
            )
            results = data.get("results", [])

            # Фильтрация по дате (инкрементальная синхронизация)
            for page in results:
                edited = page.get("last_edited_time", "")
                if edited_after and edited <= edited_after:
                    return pages   # отсортировано по убыванию — дальше только старые
                pages.append(page)

            if not data.get("has_more"):
                break
            cursor = data.get("next_cursor")

        return pages

    async def _load_page(self, page: dict) -> Document | None:
        """Получает все блоки страницы и собирает текст."""
        page_id = page["id"]
        title = _extract_rich_text(
            page.get("properties", {})
                .get("title", page.get("properties", {}).get("Name", {}))
                .get("title", [])
        )

        blocks = await self._fetch_blocks(page_id)
        lines = [_block_to_text(b) for b in blocks]
        content = "\n\n".join(l for l in lines if l.strip())

        if not content.strip():
            return None

        # URL страницы: https://notion.so/{page_id без дефисов}
        url = "https://notion.so/" + page_id.replace("-", "")

        return Document(
            page_content=content,
            metadata={
                "source":      url,
                "page_id":     page_id,
                "title":       title,
                "last_edited": page.get("last_edited_time", ""),
                "created":     page.get("created_time", ""),
                "url":         url,
                "connector":   "notion",
            }
        )

    async def _fetch_blocks(self, block_id: str, depth: int = 0) -> list[dict]:
        """Рекурсивно получаем все блоки (макс. глубина 3 уровня)."""
        if depth > 3:
            return []
        blocks = []
        cursor = None
        while True:
            url = f"{self.BASE_URL}/blocks/{block_id}/children"
            params = {"page_size": 100}
            if cursor:
                params["start_cursor"] = cursor

            data = await self._get(url, headers=self._headers, params=params)
            for block in data.get("results", []):
                blocks.append(block)
                # Рекурсивно загружаем дочерние блоки
                if block.get("has_children") and block["type"] != "child_page":
                    children = await self._fetch_blocks(block["id"], depth + 1)
                    blocks.extend(children)

            if not data.get("has_more"):
                break
            cursor = data.get("next_cursor")

        return blocks

Confluence: REST API и HTML-контент

Confluence отдаёт тело страницы в двух форматах: storage (внутренний XML) и view (HTML для браузера). Для RAG удобнее view — это стандартный HTML, который мы уже умеем парсить. Но нужно убрать навигационные элементы Confluence (панель редактирования, макросы).

# connectors/confluence_connector.py
import os
import base64
from bs4 import BeautifulSoup
from core.document import Document
from connectors.base import BaseConnector, SyncState


class ConfluenceConnector(BaseConnector):
    """
    Загружает страницы из Confluence Cloud через REST API v2.
    Аутентификация: Basic Auth (email + API token).
    """

    def __init__(
        self,
        base_url: str | None = None,       # https://your-domain.atlassian.net
        email: str | None = None,
        token: str | None = None,
        space_key: str | None = None,      # если None — загружаем всё
        max_rps: float = 5.0,
    ):
        super().__init__(max_rps=max_rps)
        self._base_url = (base_url or os.environ["CONFLUENCE_URL"]).rstrip("/")
        _email = email or os.environ["CONFLUENCE_EMAIL"]
        _token = token or os.environ["CONFLUENCE_TOKEN"]
        # Basic Auth = base64(email:token)
        credentials = base64.b64encode(f"{_email}:{_token}".encode()).decode()
        self._auth_header = f"Basic {credentials}"
        self._space_key = space_key

    @property
    def _headers(self) -> dict:
        return {
            "Authorization": self._auth_header,
            "Accept": "application/json",
        }

    async def load(self, state: SyncState | None = None) -> list[Document]:
        pages = await self._fetch_all_pages(
            modified_after=state.last_sync_at if state else ""
        )
        documents = []
        for page in pages:
            doc = await self._load_page(page)
            if doc:
                documents.append(doc)
        return documents

    async def _fetch_all_pages(self, modified_after: str = "") -> list[dict]:
        """Получаем список страниц с пагинацией (cursor-based в API v2)."""
        pages = []
        url = f"{self._base_url}/wiki/api/v2/pages"
        params: dict = {
            "limit": 250,
            "status": "current",
            "body-format": "view",   # сразу получаем HTML-тело
        }
        if self._space_key:
            params["space-key"] = self._space_key
        if modified_after:
            # Confluence v2 поддерживает фильтр по версии, но не по дате напрямую
            # Используем metadata.modified для фильтрации после загрузки
            pass

        while True:
            data = await self._get(url, headers=self._headers, params=params)
            results = data.get("results", [])
            for page in results:
                version_at = page.get("version", {}).get("createdAt", "")
                if modified_after and version_at <= modified_after:
                    continue
                pages.append(page)

            # Cursor-based пагинация через Link header
            links = data.get("_links", {})
            next_url = links.get("next")
            if not next_url:
                break
            # next_url — относительный путь
            url = f"{self._base_url}{next_url}"
            params = {}  # параметры уже в next_url

        return pages

    async def _load_page(self, page: dict) -> Document | None:
        """Парсим HTML-тело страницы и формируем Document."""
        page_id = page["id"]
        title   = page.get("title", "")
        html    = page.get("body", {}).get("view", {}).get("value", "")

        if not html:
            return None

        text = self._html_to_text(html, title)
        if not text.strip():
            return None

        space = page.get("spaceId", "")
        version = page.get("version", {})
        web_url = (
            self._base_url
            + page.get("_links", {}).get("webui", f"/wiki/spaces/{space}/pages/{page_id}")
        )

        return Document(
            page_content=text,
            metadata={
                "source":       web_url,
                "page_id":      page_id,
                "title":        title,
                "space_id":     space,
                "version":      version.get("number", 1),
                "last_modified": version.get("createdAt", ""),
                "author":       version.get("authorId", ""),
                "url":          web_url,
                "connector":    "confluence",
            }
        )

    @staticmethod
    def _html_to_text(html: str, title: str) -> str:
        """
        Конвертируем HTML Confluence в чистый текст.
        Удаляем confluence-специфичные элементы (макросы, панели).
        """
        soup = BeautifulSoup(html, "lxml")

        # Удаляем confluence-специфичные макросы и boilerplate
        for tag in soup.find_all(class_=lambda c: c and (
            "conf-macro" in c or "panel" in c or
            "aui-" in c or "confluence-information-macro" in c
        )):
            tag.decompose()

        parts = [f"# {title}"]
        for element in soup.find_all(
            ["h1", "h2", "h3", "h4", "p", "li", "td", "th", "pre", "code"]
        ):
            text = element.get_text(separator=" ", strip=True)
            if not text:
                continue
            if element.name == "h1":
                parts.append(f"# {text}")
            elif element.name == "h2":
                parts.append(f"## {text}")
            elif element.name in ("h3", "h4"):
                parts.append(f"### {text}")
            elif element.name == "li":
                parts.append(f"- {text}")
            elif element.name in ("th", "td"):
                parts.append(text)  # таблицы — упрощённо
            elif element.name in ("pre", "code"):
                parts.append(f"`{text}`")
            else:
                parts.append(text)

        return "\n\n".join(p for p in parts if p.strip())

Google Drive: Drive API + Docs API

Google Drive — это файловое хранилище, но «файлы» там бывают двух видов: обычные (PDF, DOCX, изображения) и нативные Google-форматы (Docs, Sheets, Slides). Нативные форматы нельзя скачать напрямую — их нужно экспортировать через отдельный метод.

# connectors/gdrive_connector.py
import os
import io
from pathlib import Path
from google.oauth2 import service_account
from googleapiclient.discovery import build
from googleapiclient.http import MediaIoBaseDownload

from core.document import Document
from connectors.base import SyncState
from loaders.document_dispatcher import DocumentDispatcher

# MIME-типы Google-форматов и их экспортный формат
GOOGLE_MIME_EXPORT = {
    "application/vnd.google-apps.document":     "text/plain",
    "application/vnd.google-apps.spreadsheet":  "text/csv",
    "application/vnd.google-apps.presentation": "text/plain",
}
# Обычные файлы, которые умеем читать
SUPPORTED_MIME = {
    "application/pdf",
    "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
    "text/html", "text/markdown", "text/plain",
}


class GoogleDriveConnector:
    """
    Загружает файлы из Google Drive.
    Поддерживает: Google Docs, Sheets, Slides (экспорт),
    PDF, DOCX, HTML, Markdown (прямое скачивание).
    Аутентификация: Service Account (JSON key file).
    """

    SCOPES = [
        "https://www.googleapis.com/auth/drive.readonly",
    ]

    def __init__(
        self,
        service_account_file: str | None = None,
        folder_id: str | None = None,        # None → весь Drive
        modified_after: str = "",            # ISO 8601, для инкр. синхронизации
    ):
        key_file = service_account_file or os.environ["GOOGLE_SERVICE_ACCOUNT_FILE"]
        credentials = service_account.Credentials.from_service_account_file(
            key_file, scopes=self.SCOPES
        )
        self._drive = build("drive", "v3", credentials=credentials)
        self._folder_id = folder_id
        self._modified_after = modified_after
        self._dispatcher = DocumentDispatcher()

    def load(self, state: SyncState | None = None) -> list[Document]:
        files = self._list_files(
            modified_after=state.last_sync_at if state else self._modified_after
        )
        documents = []
        for file_meta in files:
            try:
                docs = self._load_file(file_meta)
                documents.extend(docs)
            except Exception as e:
                print(f"  ✗ {file_meta['name']}: {e}")
        return documents

    def _list_files(self, modified_after: str = "") -> list[dict]:
        """Список файлов с пагинацией."""
        query_parts = ["trashed = false"]
        if self._folder_id:
            query_parts.append(f"'{self._folder_id}' in parents")
        if modified_after:
            query_parts.append(f"modifiedTime > '{modified_after}'")

        # Фильтр по поддерживаемым MIME
        supported = list(GOOGLE_MIME_EXPORT) + list(SUPPORTED_MIME)
        mime_filter = " or ".join(f"mimeType = '{m}'" for m in supported)
        query_parts.append(f"({mime_filter})")

        query = " and ".join(query_parts)
        files = []
        page_token = None

        while True:
            response = self._drive.files().list(
                q=query,
                fields="nextPageToken, files(id, name, mimeType, modifiedTime, webViewLink, parents)",
                pageSize=200,
                pageToken=page_token,
            ).execute()

            files.extend(response.get("files", []))
            page_token = response.get("nextPageToken")
            if not page_token:
                break

        return files

    def _load_file(self, file_meta: dict) -> list[Document]:
        """Скачивает файл и загружает через DocumentDispatcher."""
        file_id   = file_meta["id"]
        file_name = file_meta["name"]
        mime_type = file_meta["mimeType"]
        web_url   = file_meta.get("webViewLink", "")

        # Определяем: экспорт или прямое скачивание
        if mime_type in GOOGLE_MIME_EXPORT:
            export_mime = GOOGLE_MIME_EXPORT[mime_type]
            content = self._export_file(file_id, export_mime)
            ext = ".txt" if "text/plain" in export_mime else ".csv"
        else:
            content = self._download_file(file_id)
            ext = Path(file_name).suffix or ".bin"

        # Сохраняем во временный файл и загружаем через dispatcher
        import tempfile
        with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
            tmp.write(content)
            tmp_path = Path(tmp.name)

        try:
            docs = self._dispatcher.load(tmp_path)
        finally:
            tmp_path.unlink(missing_ok=True)

        # Обогащаем metadata Drive-специфичными полями
        for doc in docs:
            doc.metadata.update({
                "source":        web_url or f"gdrive://{file_id}",
                "file_id":       file_id,
                "filename":      file_name,
                "drive_mime":    mime_type,
                "modified_time": file_meta.get("modifiedTime", ""),
                "connector":     "gdrive",
            })

        return docs

    def _export_file(self, file_id: str, mime_type: str) -> bytes:
        request = self._drive.files().export_media(
            fileId=file_id, mimeType=mime_type
        )
        buf = io.BytesIO()
        downloader = MediaIoBaseDownload(buf, request)
        done = False
        while not done:
            _, done = downloader.next_chunk()
        return buf.getvalue()

    def _download_file(self, file_id: str) -> bytes:
        request = self._drive.files().get_media(fileId=file_id)
        buf = io.BytesIO()
        downloader = MediaIoBaseDownload(buf, request)
        done = False
        while not done:
            _, done = downloader.next_chunk()
        return buf.getvalue()

SharePoint: Microsoft Graph API

SharePoint доступен через Microsoft Graph API — единый REST-шлюз для всех Microsoft 365 сервисов. Для серверного приложения используем Client Credentials flow через MSAL (Microsoft Authentication Library).

# connectors/sharepoint_connector.py
import os
from pathlib import Path
from msal import ConfidentialClientApplication

from core.document import Document
from connectors.base import BaseConnector, SyncState
from loaders.document_dispatcher import DocumentDispatcher


class SharePointConnector(BaseConnector):
    """
    Загружает файлы из SharePoint через Microsoft Graph API.
    Аутентификация: OAuth 2.0 Client Credentials (MSAL).
    """

    GRAPH_URL = "https://graph.microsoft.com/v1.0"
    SCOPE     = ["https://graph.microsoft.com/.default"]

    def __init__(
        self,
        tenant_id: str | None = None,
        client_id: str | None = None,
        client_secret: str | None = None,
        site_id: str | None = None,          # ID сайта SharePoint
        drive_id: str | None = None,         # ID библиотеки документов
        folder_path: str = "/",              # путь внутри библиотеки
        max_rps: float = 5.0,
    ):
        super().__init__(max_rps=max_rps)
        self._site_id   = site_id   or os.environ["SHAREPOINT_SITE_ID"]
        self._drive_id  = drive_id  or os.environ.get("SHAREPOINT_DRIVE_ID", "")
        self._folder    = folder_path

        # MSAL кэширует токены автоматически
        self._msal_app = ConfidentialClientApplication(
            client_id=client_id or os.environ["AZURE_CLIENT_ID"],
            client_credential=client_secret or os.environ["AZURE_CLIENT_SECRET"],
            authority=f"https://login.microsoftonline.com/"
                      f"{tenant_id or os.environ['AZURE_TENANT_ID']}",
        )
        self._dispatcher = DocumentDispatcher()

    def _get_token(self) -> str:
        """Получаем токен из кэша или запрашиваем новый."""
        result = self._msal_app.acquire_token_silent(self.SCOPE, account=None)
        if not result:
            result = self._msal_app.acquire_token_for_client(scopes=self.SCOPE)
        if "access_token" not in result:
            raise RuntimeError(f"MSAL error: {result.get('error_description')}")
        return result["access_token"]

    @property
    def _auth_header(self) -> dict:
        return {"Authorization": f"Bearer {self._get_token()}"}

    async def load(self, state: SyncState | None = None) -> list[Document]:
        items = await self._list_items(
            modified_after=state.last_sync_at if state else ""
        )
        documents = []
        for item in items:
            try:
                docs = await self._load_item(item)
                documents.extend(docs)
            except Exception as e:
                print(f"  ✗ {item.get('name')}: {e}")
        return documents

    async def _list_items(self, modified_after: str = "") -> list[dict]:
        """Рекурсивный обход папок через Graph API."""
        items = []
        # Если drive_id не задан — берём default drive сайта
        if self._drive_id:
            base = f"{self.GRAPH_URL}/drives/{self._drive_id}"
        else:
            base = f"{self.GRAPH_URL}/sites/{self._site_id}/drive"

        path = self._folder.strip("/")
        if path:
            list_url = f"{base}/root:/{path}:/children"
        else:
            list_url = f"{base}/root/children"

        await self._collect_items(list_url, items, modified_after)
        return items

    async def _collect_items(
        self, url: str, items: list, modified_after: str
    ) -> None:
        """Рекурсивно собираем файлы, пропуская папки."""
        params = {"$top": 200, "$select": "id,name,file,folder,lastModifiedDateTime,webUrl,size"}
        if modified_after:
            params["$filter"] = f"lastModifiedDateTime gt {modified_after}"

        data = await self._get(url, headers=self._auth_header, params=params)
        for item in data.get("value", []):
            if "folder" in item:
                # Рекурсивно входим в подпапку
                folder_url = f"{url.rsplit('/children', 1)[0]}/{item['id']}/children"
                await self._collect_items(folder_url, items, modified_after)
            elif "file" in item:
                items.append(item)

        # Пагинация через @odata.nextLink
        next_link = data.get("@odata.nextLink")
        if next_link:
            await self._collect_items(next_link, items, modified_after)

    async def _load_item(self, item: dict) -> list[Document]:
        """Скачивает файл и обрабатывает через DocumentDispatcher."""
        item_id  = item["id"]
        filename = item["name"]
        web_url  = item.get("webUrl", "")

        # Ссылка на скачивание
        if self._drive_id:
            download_url_endpoint = (
                f"{self.GRAPH_URL}/drives/{self._drive_id}/items/{item_id}/content"
            )
        else:
            download_url_endpoint = (
                f"{self.GRAPH_URL}/sites/{self._site_id}/drive/items/{item_id}/content"
            )

        # Graph API возвращает 302 redirect на временную ссылку
        response = await self._client.get(
            download_url_endpoint,
            headers=self._auth_header,
            follow_redirects=True,
        )
        response.raise_for_status()
        content = response.content

        ext = Path(filename).suffix or ".bin"
        import tempfile
        with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
            tmp.write(content)
            tmp_path = Path(tmp.name)

        try:
            docs = self._dispatcher.load(tmp_path)
        except ValueError:
            return []   # неподдерживаемый тип
        finally:
            tmp_path.unlink(missing_ok=True)

        for doc in docs:
            doc.metadata.update({
                "source":        web_url,
                "item_id":       item_id,
                "filename":      filename,
                "last_modified": item.get("lastModifiedDateTime", ""),
                "file_size":     item.get("size", 0),
                "connector":     "sharepoint",
            })

        return docs

Инкрементальная синхронизация

Полная загрузка 10 000 страниц Confluence занимает несколько часов. Повторять это каждую ночь — расточительно. Инкрементальная синхронизация берёт только изменения с момента последнего запуска.

# sync/pipeline.py
import asyncio
from datetime import datetime, timezone
from pathlib import Path
from core.document import Document
from connectors.base import SyncState
from connectors.notion_connector import NotionConnector
from connectors.confluence_connector import ConfluenceConnector


STATE_DIR = Path(".sync_state")
STATE_DIR.mkdir(exist_ok=True)


async def sync_notion(full: bool = False) -> list[Document]:
    """
    Синхронизирует Notion.
    full=True → игнорировать состояние, загрузить всё заново.
    """
    state_path = STATE_DIR / "notion.json"
    state = SyncState() if full else SyncState.load(state_path)

    async with NotionConnector() as connector:
        docs = await connector.load(state=None if full else state)

    # Обновляем состояние
    now = datetime.now(timezone.utc).isoformat()
    new_state = SyncState(
        last_sync_at=now,
        processed_ids=state.processed_ids | {d.metadata["page_id"] for d in docs},
    )
    new_state.save(state_path)

    print(f"Notion: synced {len(docs)} documents (since {state.last_sync_at or 'beginning'})")
    return docs


async def sync_all(full: bool = False) -> list[Document]:
    """Параллельная синхронизация всех источников."""
    results = await asyncio.gather(
        sync_notion(full=full),
        # sync_confluence(full=full),
        # sync_gdrive(full=full),
        return_exceptions=True,
    )
    documents = []
    for r in results:
        if isinstance(r, Exception):
            print(f"Sync error: {r}")
        else:
            documents.extend(r)
    return documents


# Пример запуска
if __name__ == "__main__":
    import sys
    full_sync = "--full" in sys.argv
    docs = asyncio.run(sync_all(full=full_sync))
    print(f"Total documents ready for indexing: {len(docs)}")
💡 Инкрементальная синхронизация в продакшне требует хранить состояние в персистентном хранилище (Redis, PostgreSQL, S3), а не в локальном файле. Для удаления документов нужно отдельно отслеживать удалённые/архивированные страницы и убирать их из индекса.

Краткий обзор: что умеет каждый сервис

Notion
Block-based · JSON API
Структура: Страницы → Блоки (рекурсивно)
Поиск по дате: фильтр на клиенте (sort DESC)
Пагинация: cursor (start_cursor)
Rate limit: 3 req/s, HTTP 429 → Retry-After
Webhook: нет (нужен polling)
Сложность: рекурсия блоков, rich_text
Confluence
REST v2 · HTML body
Структура: Spaces → Pages (иерархия)
Поиск по дате: CQL: lastModified > "2024-01-01"
Пагинация: cursor через _links.next
Rate limit: 10 req/s для Cloud
Webhook: есть (Atlassian Connect)
Сложность: HTML-макросы, body formats
Google Drive
Drive API v3 · Service Account
Структура: Files (плоский + папки)
Поиск по дате: modifiedTime > '…' в query
Пагинация: pageToken
Rate limit: 1000 req/100s/user
Webhook: есть (Push notifications)
Сложность: нативные форматы → export
SharePoint
Graph API · MSAL OAuth
Структура: Sites → Drives → Items
Поиск по дате: $filter=lastModifiedDateTime gt …
Пагинация: @odata.nextLink
Rate limit: 10 000 req/10min
Webhook: есть (change notifications)
Сложность: MSAL token refresh, Azure AD

Типичные ошибки

Не кэшировать OAuth-токен
Запрашивать новый access token перед каждым API-вызовом — это 1 лишний HTTP-запрос на каждый вызов и риск превысить лимит запросов к auth-серверу. SharePoint токены живут 3600 секунд.
→ MSAL кэширует автоматически. Для других OAuth — сохраняйте токен и expires_at, проверяйте перед запросом.
Не обрабатывать HTTP 429
429 Too Many Requests — штатная ситуация при больших объёмах. Без обработки коннектор упадёт на 500-й странице и придётся начинать сначала без сохранения прогресса.
→ Retry с exponential backoff + Retry-After заголовок + сохранение cursor в SyncState.
Загружать все поля Notion-страниц
Notion Search API возвращает все свойства страниц по умолчанию. Для баз данных с сотнями кастомных полей это огромный JSON. Плюс блоки загружаются отдельными запросами — без ограничения глубины рекурсии можно уйти в бесконечный обход.
→ Ограничивайте глубину рекурсии (depth ≤ 3). Запрашивайте только нужные поля через filter_properties.
Хранить секреты в коде
API токены и client_secret в репозитории — критическая уязвимость. Даже в приватном репозитории токены нужно ротировать после случайного коммита.
→ Только environment variables или secret manager (Vault, AWS Secrets Manager). python-dotenv для локальной разработки.
Полная перезагрузка вместо инкрементальной
Каждую ночь грузить 50 000 страниц Confluence целиком — нормально для MVP, катастрофа для production. При сбое теряется весь прогресс.
→ С первого дня сохраняйте cursor и last_sync_at. Incremental sync — это не оптимизация, это надёжность.

Шпаргалка

📋 Выбор метода аутентификации и ключевые паттерны
  • Notion — API key (secret_…), 3 req/s, блоки → рекурсия, Search API для листинга.
  • Confluence — Basic Auth (email:API token в base64), REST v2, body-format=view → HTML → BeautifulSoup.
  • Google Drive — Service Account JSON, нативные форматы → export_media(), обычные → get_media().
  • SharePoint — MSAL Client Credentials, Graph API, @odata.nextLink для пагинации, 302 redirect на download.
  • Rate limiting — токен-баккет перед каждым запросом, Retry-After при 429.
  • Инкрементальная синхронизация — сохранять last_sync_at и cursor, фильтровать по modified date.
  • Секреты — только env vars, никогда в коде.
  • Результат — всегда list[Document] с source, connector, last_modified в metadata.

Практическое задание

  1. Notion коннектор. Создайте интеграцию в Notion, расшарьте на неё 2–3 страницы. Реализуйте NotionConnector и загрузите страницы. Убедитесь, что в metadata есть title, page_id и last_edited. Попробуйте страницу с вложенными блоками — проверьте глубину рекурсии.
  2. Rate limiter под нагрузкой. Напишите тест: создайте RateLimiter(max_rps=5) и выполните 20 await limiter.acquire() подряд, замерив общее время. Убедитесь, что 20 вызовов занимают ≥ 4 секунды (20 / 5 rps = 4s). Поэкспериментируйте с параметром.
  3. Инкрементальная синхронизация. Запустите загрузку Notion дважды: первый раз — полная загрузка, второй — инкрементальная (с сохранённым SyncState). Измените одну страницу в Notion между запусками. Убедитесь, что второй запуск загружает только изменённую страницу, а не все.