Слой ingestion — как это работает на самом деле

Статус: current (as-built), сверено на commit 544d4e51. Целевая сервисная граница ingestion описана в ../../30-services/ingestion/README.md. Этот документ — не альтернатива ей, а проверка по коду: что из целевого дизайна реально существует, как называются реальные символы и где проходит реальный путь выполнения.

Назначение

Слой ingestion отвечает за то, чтобы данные от поставщиков (цены, остатки, характеристики) регулярно попадали в систему и доезжали до нормализации. В коде это не отдельный сервис, а worker-процесс supplier-sync, который поднимается как fx-приложение и работает по тикерам — без публичного HTTP.

Участники

  • supplier-sync (backend/cmd/supplier-sync/main.go) — процесс-контейнер: поднимает fx-граф, запускает per-supplier тикеры, отдаёт loopback /healthz, /readyz, /metrics. Точка входа зафиксирована в ../entry-points.md.
  • Pipelines (backend/internal/core/ingestion/app/) — PipelineRouter выбирает пайплайн по Capabilities поставщика: BulkSnapshotPipelineImpl для batch-поставщиков (Russvet, ETM, IEK, Systeme) или cursor-пайплайн для DKC.
  • Connectors (backend/internal/core/ingestion/infra/<supplier>/ и смежные пакеты) — HTTP-клиенты конкретных поставщиков, вызываемые из пайплайнов.
  • Normalization — преобразование сырых ответов поставщика во внутреннюю модель (NormalizedOffer, Observation); подробный as-built разбор появится отдельным use-case файлом.
  • Offers — запись нормализованных данных в PostgreSQL (supplier_offers, offer_observations) и транзакционный outbox для Kafka.

Как это работает: тик и диспетчер конвейеров

Тик (tick) — одно «срабатывание часов»: момент, когда воркер supplier-sync просыпается по таймеру и делает одну порцию работы для одного поставщика. В проде на каждого поставщика висит свой тикер со своим интервалом (напр. russvet — раз в 30 минут / 12 часов). Каждое срабатывание = один тик = «сходи к поставщику и забери данные один раз». В тестах тик дёргают вручную, без таймера.

PipelineRouter (backend/internal/core/ingestion/app/pipeline_router.go) — диспетчер. Сам данные не тянет: на каждый тик он смотрит «возможности» поставщика (Capabilities, паспорт «как этот поставщик умеет отдавать данные») и передаёт работу одному из трёх конвейеров:

КонвейерКогда выбираетсяЧто делаетПример
cursorIncrementalSync = revision/dateтянет только изменения с прошлой меткиDKC
bulkScheduledRefresh.UsesBulkSnapshot() + есть источниккаждый тик выгружает каталог пачкойRussvet, ETM, IEK, Systeme
fullиначе (fallback)legacy полный каталог

Порядок проверки в RunTick: cursor → bulk → full. Для Russvet выбирается bulk. Диспетчер также пишет метрику «какой конвейер выбран» (PipelineRouted) — это телеметрия, на маршрутизацию она не влияет.

Отдельно: у Russvet сам bulk-конвейер разделён на два job kindcatalog (описания + характеристики) и commerce (цены + остатки); дефолтный тик обрабатывает catalog, коммерческие оси приходят отдельным commerce-тиком.

Use-case файлы слоя

  • russvet.md — сквозной путь Russvet: supplier-sync → БД/raw-хранилище → charnorm-worker → matcher-worker/LLM → canonical-assignment-worker → БД. Появится следующей задачей инициативы.

По мере разбора добавляются use-case файлы для остальных поставщиков и внутренних переходов слоя (нормализация, запись офферов).

Куда дальше