← Harvester

Модель данных

harvester · workzone

Страница о том, какие данные течёт сквозь модуль и как они живут во времени. Два контракта несут конвейер — RawItem Entity — затем Entity расходится по хранилищам, а четыре решения задают её судьбу после загрузки: что мы помним из истории, как обходимся с удалениями и вложениями, и где проходит граница с Knowledge Store. Помимо этого контент-плана модуль держит свой контрольный слой — конфигурацию источников и журнал прогонов; его правит через API Admin Panel.

RawItem
Connector → Transform

Сырьё источника как есть, вместе с его метаданными. Выход коннектора и вход трансформации: ещё не нормализован, формат каждого источника свой.

  • полезная нагрузка сырой ответ источника — тело тикета, страницы, сообщения как их отдал API
  • идентификаторы тип и id записи у источника — будущий ключ идемпотентности
  • метаданные сбора когда и каким прогоном получено, курсор инкремента, отметка источника
Entity
Transform → Loader

Нормализованная сущность — единая форма для всего, что собирает модуль. Разные типы внутри источника (тикет, страница, сообщение) — подвиды одной Entity и проходят один конвейер.

  • идентичность ключ идемпотентности source_id + source_type + source_entity_id — по нему Loader делает upsert: ровно одна строка на запись источника
  • тип подвид сущности внутри источника (ticket · page · message)
  • статус draft · final · archived — жизненный цикл контента в источнике, выставляет шаг classify. Пропажа записи из источника — отдельная ось (удаления), не значение статуса
  • контент текст для выдачи и его эмбеддинги — по фрагментам (chunks), не один вектор на запись
  • метаданные автор, временные метки, ссылки внутри источника и упоминания чужих источников (refs)
  • ACL права доступа источника, перенесённые вместе с записью — Harvester их захватывает, а модель и хранение держит Knowledge Store
Куда ложится Entity: три проекции хранения

Один логический контракт Entity на загрузке расходится по трём парадигмам хранения — это не одна запись в одной базе. Точную раскладку (какое поле в реляционную, векторную или графовую базу) держит → Knowledge Store; здесь — предварительная картина, уточняемая по факту проектирования хранилища.

Entity
Реляционная структура · точные фильтры

Тело записи — идентичность, тип, статус, метаданные, текст для выдачи — и ACL как source of truth прав. Всё, что фильтруют и джойнят точно.

Векторная смысл · поиск по близости

Эмбеддинги фрагментов (chunks) — дочерние строки сущности (ключ — её id) с контент-хешем: опора семантического поиска по смыслу, а не по точному совпадению слов.

Графовая связи · обход

Сущность как узел и связи внутри источника, что строит конвейер. Упоминания чужих источников (refs) станут рёбрами позже — уже на стороне Knowledge Store.

Права на выдаче применяются фильтром: фрагмент наследует ACL своей сущности (ключ — её id), а ACL лежит в реляционной базе рядом с векторами — поэтому отбор по правам идёт тем же запросом, что и поиск по смыслу. Захват прав и членства держит acl-identity, а их модель и хранение — Knowledge Store.

История: что мы помним во времени

Три разных вопроса об изменениях, и на каждый — свой ответ. Их легко смешать, поэтому решение разнесено явно: одно храним сейчас, другое откладываем, третье не храним сознательно.

снимок
Сущность — текущий снимок
Тикет, страница, сообщение: храним последнее состояние. Повторный прогон перезаписывает строку через upsert; прошлых редакций самой сущности не накапливаем.
накопление v2
Обсуждения — накопительно
Комментарии и треды — отдельные накопительные единицы рядом с сущностью, а не перезаписываемое поле. Тред растёт со временем, каждая реплика ценна сама по себе. Отложено в v2.
не храним
Changelog полей — не храним
Что поле было Open и стало Done, по шагам — не ведём. Только updated_at отмечает факт изменения; пожизненную ленту правок каждого поля не реконструируем.
Удаления
soft-delete · отдельная ось

Запись, исчезнувшую в источнике, помечаем удалённой и скрываем из выдачи, но физически храним — это даёт аудит (что было и когда ушло) и восстановление, если запись вернулась или удаление было ошибкой. Маркер удаления — отдельная ось от статуса, не его значение: пропажу из источника видит полный скан reconciliation (в инкремент удаление не приходит) — он и проставляет момент исчезновения. Статус классификации (draft · final · archived) при этом сохраняется как был, поэтому восстановление однозначно: снимаем маркер — прежний статус возвращается сам. archived остаётся за контентом (автор заархивировал — запись в источнике есть), удаление — за присутствием; держать оба смысла в одном значении нельзя.

Вложения
Текст

В v1 берём только то, что уже текст: тело письма, текстовый файл, сообщение с вложением-текстом. Это сразу годится для индексации и эмбеддинга.

Парсинг файлов v2

Извлечение текста из PDF и docx, распознавание сканов (OCR) — отдельный слой обработки. Отложено в v2: бинарные вложения пока сохраняются как ссылки, без разбора содержимого.

Граница cross-source с Knowledge Store

Связи и дедупликацию делят два модуля. Правило границы простое: что можно установить точно и внутри одного источника — делает Harvester; что требует догадки между источниками — отдаёт → Knowledge Store.

Harvester точно · внутри источника
  • Связи внутри источника — родитель-потомок, ссылки тикета на тикет того же трекера — строит сам, прямо в конвейере.
  • Ссылки на чужие источники — сохраняет заявкой (refs), не материализует в рёбра: целевой узел может быть ещё не загружен. Формат — ниже.
  • Лёгкая дедупликация — exact match по source_id + source_type + source_entity_id: повторный прогон узнаёт ту же запись и обновляет её.
  • Identity по email — сведение учёток по точному совпадению адреса; явный ключ держит acl-identity.
Knowledge Store догадка · между источниками
  • Cross-source рёбра — связывает узлы из разных источников по явным упоминаниям, когда оба уже в графе.
  • Нечёткое entity resolution — сводит сущности по похожести (fuzzy matching), а не по точному ключу.
  • Глубокая дедупликация — то, что exact match не поймал: один объект под разными обличьями в разных системах.
  • «Разные email, тот же человек» — нечёткое сведение личностей; либо entity resolution здесь, либо ручной разбор в Admin.
Формат refs таблица entity_ref · Knowledge Store · до материализации

refs живут отдельной строкой в таблице entity_ref (дом — Knowledge Store) — заявка на связь, ещё не ребро. Заполняет нормализация; каждая несёт минимум, достаточный, чтобы → Knowledge Store позже разрешил цель и материализовал ребро в entity_edge, даже если целевой узел на момент захвата ещё не загружен.

  • relation тип связи · mentions · blocks · parent · author · duplicate · …
  • target_kind вид цели · issue · page · user · commit · mr · message
  • target_ref натуральный идентификатор как он есть в источнике · ключ "PROJ-123" · URL · "@user" · sha
  • source_hint предполагаемый источник / коннектор · jira · gitlab · … · null, если неизвестен
Источник, прогоны, сбойные items

Три конкретные таблицы, дом которых — Harvester: подключённый источник, журнал его прогонов и очередь сбойных items на разбор. Экраны Admin Panel (хаб и карточка источника) читают и правят их через API — поведением владеет модуль.

1 Источник Конфигурация подключения.
sources подключённые источники
секреты at-rest
id BigInteger PK
name Text NOT NULL отображаемое имя · «Jira · Acme»
connector_type Text NOT NULL ключ коннектора · набор открытый, без CHECK · → манифест
base_url Text NULL адрес инстанса · у части типов не нужен
auth_account Text NOT NULLDEFAULTCHECK под кем · service · personal
auth_method Text NOT NULLDEFAULTCHECK чем · static_token · oauth v2
credential_enc Text NULL шифротекст AES-256-GCM · NULL у Disconnected · → секреты
state Text NOT NULLDEFAULTCHECK намерение админа · active · paused · disconnected
scope_mode Text NOT NULLDEFAULTCHECK всё / только выбранное
scope_list JSONB NOT NULLDEFAULT deny-list (all) или allow-list (selected): id контейнеров
content_filters JSONB NOT NULLDEFAULT тумблеры «сверх правил» · набор задаёт манифест
sync_interval Integer NULL сек · интервал инкремента · override глоб. default · NULL = наследует
reconcile_interval Integer NULL сек · каденс reconciliation · override глоб. default · NULL = наследует
reconcile_window Integer NULL минута недели старта reconciliation, org-tz · суточная каденс → минута суток (mod 1440) · override глоб. default · NULL = наследует
authority_tier Text NULLCHECK уровень доверия источнику · low · normal · high · вход KS decay · числовой множитель выводит код · NULL = default_authority манифеста по типу
webhook_enabled Boolean NOT NULLDEFAULT канал реал-тайм · DEFAULT false
webhook_secret_enc Text NULL шифротекст HMAC-секрета · → секреты
incremental_cursor JSONB NULL позиция since прошлого прогона · NULL до первого
last_probe_at DateTime(tz) NULL время плановой пробы соединения · NULL до первой
last_probe_status Text NULLCHECK итог пробы · ok · unreachable · auth_failed
created_at DateTime(tz) DEFAULT server_default=now()
updated_at DateTime(tz) DEFAULT now() + trigger
Состояние хранится, здоровье вычисляется → Жизненный цикл источника
В БД лежит state — намерение админа (active · paused · disconnected). Здоровье (idle · syncing · error) не отдельная колонка: его выводят из последнего прогона и плановой пробы соединения. Сам итог пробы персистится (last_probe_at · last_probe_status) — это единственная грань здоровья, которую из прогона не вывести (лёгкая проба между прогонами падает отдельно). Бейджи в хабе — отрисовка двух осей, не третья истина.
Scope и фильтры — политика в JSONB → Манифест коннектора
scope_mode + scope_list и content_filters — гибкая политика, не замороженный снимок: какие объекты scope и какие тумблеры «сверх правил» доступны, объявляет манифест коннектора, поэтому форма колонок не фиксирована — отсюда JSONB.
Секреты — крипто-ядро Auth → Крипто-ядро
credential_enc и webhook_secret_enc — шифротекст AES-256-GCM; шифр и ключ держит крипто-ядро Auth, своего хранилища Harvester не заводит. Наружу (API, экспорт) значения не уходят — UI получает лишь маску.
Глобальный default расписания
Каденс-настройки сбора — sync_interval (интервал инкремента), reconcile_interval (частота сверки) и reconcile_window (когда в неделе её запускать) — на источнике переопределяют глобальные значения по умолчанию (NULL = наследует). Дом самих default — platform_settings (core): один глобальный набор на платформу, поле источника лишь точечно его перебивает. Окно — у тяжёлой редкой сверки (ночь / выходной); у интервального инкремента его нет. Минута хранится смещением без зоны — час admin задаёт в таймзоне организации (platform_settings.timezone), планировщик разворачивает к ближайшему UTC-моменту запуска. Глобальный дефолт окна — воскресенье 03:00.
Реал-тайм без своих колонок: дедуп и watchdog → Webhook'и
Защита от повторной доставки (один webhook — один раз) не нуждается в колонке: виденные delivery-id живут эфемерной TTL-памятью в Cache & Workers (Redis), окно дедупа эфемерно. Точку отсчёта watchdog (молчание webhook-источника дольше окна) тоже не храним отдельно — её даёт время последнего прогона источника.
2 Прогоны синхронизации Журнал прогонов. → SyncRun
sync_runs журнал прогонов
FK → sources один активный на источник
id BigInteger PK
source_id BigInteger FK→sourcesIDX CASCADE
mode Text NOT NULLCHECK стратегия выборки · full · incremental · reconciliation
trigger Text NOT NULLCHECK что инициировало · connect · schedule · webhook · watchdog · manual
state Text NOT NULLDEFAULTCHECK queued · running · succeeded · failed · cancelled
scope JSONB NULL подмножество прогона · NULL = весь источник · окно · список items (DLQ) · контейнеры
entities_done Integer NULL обработано (N) · итог при succeeded
entities_total Integer NULL оценка объёма (M) для прогресса «N из M»
checkpoint JSONB NULL позиция возобновления внутри прогона
error_count Integer NOT NULLDEFAULT число в DLQ · DEFAULT 0
error_detail Text NULL краткая причина провала
started_at DateTime(tz) NULL NULL пока queued · длительность = finished − started
finished_at DateTime(tz) NULL терминальное время
heartbeat_at DateTime(tz) NULL живость воркера
created_at DateTime(tz) DEFAULT now()
Один активный прогон на источник
Замок — partial unique index UNIQUE (source_id) WHERE state IN ('queued','running'): один источник несёт максимум один незавершённый прогон, на время которого карточка под замком. Cancel терминализует прогон (cancelled) — следующий стартует свежим.
Частичный успех — это succeeded с ошибками
Сбойные items не валят прогон: при них он терминируется как succeeded с error_count > 0 — данные дошли, часть отложена в DLQ на разбор. failed — это полный провал прогона (источник недоступен, прогон сорван целиком). Поэтому state и error_count читаются вместе: зелёный с пометкой об ошибках — не красный.
Три оси прогона: стратегия · триггер · scope → Режимы синхронизации
Прогон описывают три независимых поля, не одно. mode — стратегия выборки (full · incremental · reconciliation, ровно три режима sync-modes): окно since и агрессивность resolve. trigger — что подняло прогон. scope — какое подмножество тянем. Все три живут у прогона, не у источника: источник держит лишь incremental_cursor и расписание. Так partial re-sync и dlq retry — не отдельные режимы, а incremental с разными триггером и scope: первый — watchdog + окно восстановления, второй — manual + список items из DLQ. «Тип» строки в истории выводится из этой тройки — отдельной колонкой не хранится. Провал прогона поднимает уведомление тем же каналом, что и прочие алерты платформы.
Что запустило — здесь, кто запустил — в журнале платформы → Audit Log
trigger несёт что подняло прогон (connect · schedule · webhook · watchdog · manual) — видно в истории прогонов, дополняет mode. А кто из админов запустил ручной прогон или менял конфигурацию источника, Harvester у себя не дублирует: личность актора, действие и результат пишет платформенный append-only журнал админ-действий — его дом Auth & Security. Аудит живёт в одном месте, а не размазан по таблицам модулей.
Метрики прогона vs объём источника
Числа прогона — его собственные: обработано (entities_done из оценки entities_total), «N из M» в прогрессе. А итоговый объём источника (колонка «Сущностей» в хабе) — не счётчик Harvester, а COUNT узлов источника в Knowledge Store; своей колонки под него модуль не заводит.
3 Сбойные items Очередь разбора. → DLQ
dead_letters очередь сбойных items
FK → sources FK → sync_runs одна строка на item
id BigInteger PK
source_id BigInteger FK→sourcesIDX CASCADE · очередь живёт, пока жив источник
run_id BigInteger FK→sync_runsNULL прогон последнего падения · SET NULL переживает чистку журнала
source_type Text NOT NULL тип записи у источника (ticket · page · message)
source_entity_id Text NOT NULL id записи у источника · что перетягивает повтор
reason Text NOT NULLCHECK категория для разбора · permission · not_found · malformed · rate_limited · unknown
error_detail Text NULL текст ошибки · для диагностики
attempts Integer NOT NULLDEFAULT сколько раз падал · растёт при повторном попадании
created_at DateTime(tz) DEFAULT впервые отложен
updated_at DateTime(tz) DEFAULT последнее падение · now() + trigger
Одна строка на сбойный item
Ключ дедупа — UNIQUE (source_id, source_type, source_entity_id): тот же item, упавший снова в следующем прогоне, не плодит дубль, а обновляет строку (attempts, updated_at, свежая reason) через upsert. Очередь — снимок актуально-нерешённого, а не лента всех падений.
Transient или permanent — решает коннектор → Ретраи
Категория reason кодирует разбор, но попадёт ли item в очередь, зависит от другой оси — преходящая ошибка (transient) или постоянная (permanent). Эту классификацию объявляет коннектор: дефолтный HTTP-маппинг (429 · 5xx → transient; 403 · 404 → permanent) он может переопределить под специфику своего API. Саму механику попыток и backoff держит reliability — сюда же приходит лишь то, что ретраи исчерпали.
Уход из очереди — удаление при успехе → Видимость сбоев
«Повторить сбойные» запускает incremental-прогон вручную (trigger=manual) со scope из нерешённых строк источника — это и есть dlq retry. Item обработан — строка удаляется (его актуальное состояние теперь в Entity через upsert); снова упал — остаётся и обновляется; исчез в источнике — удаляется, разбирать нечего. Это рабочая очередь, не журнал: аудит «сколько падало и когда» уже несёт sync_runs.error_count, поэтому решённое чистим без дублирующей истории.
Счётчик прогона ≠ глубина очереди → Разбор сбойных
sync_runs.error_count — неизменный снимок: сколько упало в том прогоне. Счётчик «K в DLQ» в хабе и разбор по причине — это живой dead_letters источника (COUNT и GROUP BY reason); по мере починки он тает, тогда как error_count прогона остаётся как был. Набор категорий reason — source of truth бэкенда, UI лишь рисует.