Russvet — путь ingestion (as-built)

Статус: current (as-built), сверено на commit 544d4e51. Целевой дизайн: ../../30-services/ingestion/ и ../../10-business/scenarios/ingestion-flow.md. Этот документ описывает, как поток реально исполняется в коде; расхождение с target-документом — повод для ADR/обновления target, а не для правки здесь.

Назначение

Russvet — 5-й поставщик (REST/JSON, Basic auth, RUB-only, ETIM). Поставщик batch-first: per-SKU Connector.Fetch не поддержан (errFetchNotSupportedconnector.go:26), весь сбор идёт через BulkSnapshotSourceBulkSnapshotPipeline. Собираются 4 оси данных: описания, характеристики, цены, остатки.

Точка входа

runSupplierTick(ctx, "russvet", ...)backend/cmd/supplier-sync/main.go:380, вызывается по тикеру, зарегистрированному в registerPerSupplierTickersmain.go:306. Публичного HTTP у процесса нет (см. ../entry-points.md, раздел «Явные отсутствия»).

Предусловия

  • В конфиге задан SUPPLIER_RUSSVET_BASE_URL, а SUPPLIER_SYNC_SUPPLIERS содержит russvet:...:true.
  • Учётные данные Basic: SUPPLIER_RUSSVET_LOGIN/SUPPLIER_RUSSVET_PASSWORD как fallback либо per-credential из secrets-пула (Credentials BC).
  • В БД желателен canonical_id для привязки observation-событий проектора: без него проекционные события (price.observed/stock.observed) не эмитятся, но supplier_offers + offer_observations + outbox_events всё равно пишутся.

Sequence-диаграмма

sequenceDiagram
    autonumber
    participant TICK as supplier-sync ticker
    participant TS as runSupplierTick (main.go:380)
    participant LG as LeaseGuard (app/lease.go)
    participant PR as PipelineRouter (pipeline_router.go:68)
    participant BP as BulkSnapshotPipeline (bulk_snapshot_pipeline.go:223)
    participant BS as russvet.BulkSnapshotSource (bulk_snapshot_source.go)
    participant API as 🌐 Russvet API
    participant NP as Pipeline.Normalize (normalization/app/pipeline.go:39)
    participant RN as russvet.Normalizer (normalization/infra/russvet)
    participant US as UpsertService.Upsert (offers/app/upsert_service.go:67)
    participant DB as PostgreSQL
    participant S3 as MinIO (raw)

    TICK->>TS: тик по интервалу
    TS->>LG: Acquire("ingestion:russvet")
    LG-->>TS: release() | ErrLeaseBusy
    TS->>PR: RunTick(ctx, "russvet")
    Note over PR: Capabilities.ScheduledRefresh.UsesBulkSnapshot()==true
    PR->>BP: RunTick(ctx, russvet)
    BP->>BS: StreamRequest() (bulk_snapshot_source.go:340)
    BS->>API: Basic auth session (session.go)
    BS->>API: GET /stocks (fetchStocks :1360)
    BS->>API: GET /position/{stock}/all (fetchCatalog :1372, paged)
    BS->>API: GET /residue/all/{stock} (fetchResiduePage :2035)
    BS->>API: GET /partnerwhstock/all/{stock} (fetchPartnerStockPage :2085)
    BS->>API: POST /massprice (fetchPrices :1455, батчи по 50)
    BS->>API: GET /specs/{sku} (fetchSpecs :2137, throttled)
    API-->>BS: goods / price / remains payloads
    loop на каждый SKU-результат
        BP->>S3: Put(raw goods/price/remains)
        BP->>NP: Normalize(ctx, "russvet", FetchResult{KindGoods,KindPrice,KindRemains})
        NP->>RN: Normalize(...)
        Note over RN: MapOffer (описания+категории) offer_mapper.go:92<br/>collectCharacteristics (ETIM specs) offer_mapper.go:213<br/>MapPrice (RUB set) price_mapper.go:28<br/>MapStock (residues+partner pool) stock_mapper.go:55
        RN-->>NP: NormalizedOffer + Observation
        NP-->>BP: NormalizedOffer + Observation
        BP->>US: Upsert(ctx, NormalizedOffer, Observation)
        US->>DB: BEGIN
        US->>DB: OfferRepository.InsertOrTouch → supplier_offers (offer_repo.go:263)
        US->>DB: ObservationRepository.Append → offer_observations (observation_repo.go:158)
        US->>DB: outbox.Repository.Append → outbox_events "offer.observation.v2" (repo.go:37)
        US->>DB: COMMIT
    end
    BP-->>TS: outcome
    TS->>TS: метрика TickAttempt(russvet, outcome)

Маппинг 4 осей

Нормализатор russvet.Normalizer.Normalize (normalizer.go:22) связывает мапперы в NormalizedOffer + Observation.

ОсьСимвол (file:line)Как маппится
ОписанияMapOfferoffer_mapper.go:92Name: specs.INFO[0].DESCRIPTION → fallback position.NAME; Description: LONG_DESCRIPTION → fallback DESCRIPTION; бренд: specs.INFO[0].BRAND → fallback position.BRAND; категории — из RS_CATALOG (Level2/3/4), fallback ETIM-класс (ETIM_CLASS/ETIM_CLASS_NAME). pos.Category («СКЛ»/«ЗКЗ») в таксономию НЕ попадает — это статус наличия.
ХарактеристикиcollectCharacteristicsoffer_mapper.go:213ETIM SPECS[]CharacteristicRaw{SupplierCode=FEATURE_CODE, Value, UOM, ...}. Не персистятся в supplier_offers — уходят вниз в charnorm.
ЦеныMapPriceprice_mapper.go:28Валюта RUB (helpers.go:20). Net=Personal, Gross=Personal_w_VAT, ListTarif=MRC_w_VAT (fallback Retail), RetailRec=Retail_w_VAT. Если все пусты — PriceSet.OnRequest=true.
ОстаткиMapStockstock_mapper.go:55residues[]StockCurrent["regional_center:{warehouse.id}"] c QtyRC=RESIDUE (собственный on-hand, per-организация). Partner pool (residue.partnerQuantityInfo + partnerWarehouseStock, по всем организациям) схлопывается в одну StockForecast.Incoming запись: Qty=max (единый общий пул, НЕ сумма) + ETA=ранняя, склад manufacturer_warehouse:manufacturer.

Побочные эффекты

  • БД (одна tx в UpsertService.Upsert): upsert supplier_offers; append offer_observations; insert outbox_events (топик offer.observation.v2). При заданном canonical_id — проекционные события price.observed/stock.observed в catalog_projector_queue.
  • Внешние вызовы: Russvet REST — /stocks, /position/.../all, /residue/all/..., /partnerwhstock/all/..., /massprice, /specs/....
  • Хранилище: raw payloads в MinIO (ссылки в offer_observations.raw_refs).
  • События: Kafka offer.observation.v2 через транзакционный outbox-диспетчер.

Выход

Актуальный SupplierOffer + новое Observation в БД. Downstream: charnorm-worker (нормализация характеристик по событию) и matcher-worker (сопоставление с canonical).

Неоптимальные места (бэклог — не чинить в рамках этой итерации)

  • bulk_snapshot_source.go — 2683 строки: в одном файле смешаны транспорт, пагинация, батчинг, бюджеты тиков. Кандидат на разбиение по осям.
  • Два пути стрима сосуществуют: StreamRequest (cursor-aware, :340) и legacy streamFull (:179) — двойная логика, риск расхождения поведения.
  • russvet/connector.goConnector существует только как заглушка (FetcherrFetchNotSupported, :26); интерфейс ingdom.Connector реализуется формально. Кандидат на отдельную capability-модель для batch-only поставщиков.
  • offers/app/upsert_service_integration_test.go:94 и :163 вызывают NewUpsertService с 5 аргументами при текущей 6-арг сигнатуре (добавлен deriverupsert_service.go:40). Тест под -tags=integration не компилируется; build-tag не покрыт обычным lint.
  • Миграция migrations/0234_canonical_analog_assignment_lookup.sql объявляет $$-функции без директив -- +goose StatementBegin/StatementEnd — встроенный goose (pgstorage.RunMigrations) рвёт тело функции по первому ; («unterminated dollar-quoted string»). Все интеграционные тесты через pgtest, чья схема доходит до 0234, ломаются на применении миграций. psql выполняет файл корректно (директивы — комментарии), поэтому прод/CLI-путь не задет; кандидат — добавить аннотации в 0234.
  • Нестабильный порядок страниц /position/{wh}/all: между сессиями russvet отдаёт каталог в разном порядке. Независимые catalog- и commerce-свипы при ограничении по объёму хватают непересекающиеся SKU. Для адресной догрузки цен/остатков к известным SKU добавлен streamCommerceTargeted (см. ниже).

Контур и тест

Russvet ingestion разделён на два job kind: catalog (JobKindCatalog, дефолт RunTick) — описания + характеристики (goods + /specs); commerce (JobKindCommerce) — цены (/massprice) + остатки (/residue, /partnerwhstock). Один дефолтный тик наполняет только описания/характеристики; цены/остатки приходят отдельным commerce-тиком.

Полный живой контур покрыт тестом backend/internal/tests/russvet_contour/russvet_live_contour_test.go::TestRussvetLiveContour_AllAxesPopulated (build tag contour). Он прогоняет в одном процессе, синхронно: catalog-тик → commerce-тик (точечно по SKU из БД, streamCommerceTargeted) → charnorm.RunBatch (LLM, с WithOutbox → эмит MappingResolved) → facts-projector → matcher.RunBatch (LLM Tier-3) → canonical-assignment.RunBatch, и проверяет наполненность набора карточек по осям (таблица «SKU × ось»).

Требует поднятый local-prod стек (Postgres/Kafka/MinIO), доступ к живому API russvet, боевой russvet-кред в пуле (Credentials BC) и живой LLM-gateway. Запуск:

cd backend && go test -tags=contour ./internal/tests/russvet_contour -run TestRussvetLiveContour -v -timeout 25m
# под дебагом:
dlv test -tags=contour ./internal/tests/russvet_contour -- -test.run TestRussvetLiveContour -test.v

Ключевой инвариант, который тест защищает: charnorm обязан эмитить MappingResolved в outbox (matching.char_facts.v1) — иначе facts-projector не строит offer_characteristic_facts, а canonical-assignment (dirty-from-facts) не обогащает canonical. Пустые offer_characteristic_facts при непустых char_name_mappings — маркер отвалившегося outbox-эмита.