План stage4 pipeline
stepan edited this page 2026-09-13 00:17:00 +03:00
This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

Перенесено из репозитория (docs/superpowers/plans/2026-09-05-deal-stage4-pipeline.md). Актуальная версия — здесь, в вики.

Дейл (Deal) — Этап 4: Pipeline и «Обработка»: очередь, стоп-лист, дедуп, отсев, ML/ИИ-порты, FTS Implementation Plan

Исторический документ этапа 4. Актуальное состояние — docs/superpowers/STATUS.md и docs/technical/Техническая-документация-Дейл.md.

Goal: Оживить в модульном монолите src/core вкладку «Обработка» Vue-фронта 1:1-контрактом /api пайплайна входящих: приём сообщений (порт + демо-ингвест до telegram-этапа 6), очередь сырых сообщений, разбор фоновым воркером по пути ТЗ §5 источник → очередь → стоп-лист → дедуп → ML → ИИ → карточка, отсев с причиной/источником решения (правила/ML/ИИ/система + конкретное слово/фраза), возврат из отсева (ignore-причин + обучение), полнотекстовый поиск по отсеву и карточкам (настоящий FTS в Postgres), автоочистка отсева раз в 3 суток + ручная, счётчики вкладки. К концу этапа ProcessingView полностью обслуживается бэкендом на реальном сквозном пути «демо-сообщение → очередь → фильтры → карточка/отсев» (telegram-источник — этап 6); приёмка — unit/curl/psql. ML-модель не готова (LocalMlClient ready:false) — ML-ветка реализована, но «спит» до этапа 6; ИИ — порт IAiClassifier + детерминированный локальный классификатор (реальный ai-service — этап 6).

Architecture: новый модуль Deal.Modules.Pipeline (чистый, без EF/HTTP) — владелец таблиц QueueItems/RejectedItems/DedupEntries (миграция TenantPipeline в TenantDbContext) и логики воркера: DTO очереди/отсева (§4.5), порт IPipelineStore, сервисы PipelineIngestService (приём, используется демо-ингвестом и, на этапе 6, gRPC-адаптером telegram-service), PipelineProcessingService (чтение/поиск очереди и отсева, возврат, очистки, запись отсева), чистое ядро разбора MessageParseCore (clean_short/clean_block, normalize_list/stack, qualify/build/primary контакты, dedup-хэш, compose_summary «О заявке», локальные поля _local_fields), PipelineWorkerService (pump: stale → stage1 → дедуп → ML → ИИ/локальный разбор → карточка/отсев). Настройки — порт ISettingsStore + IncomingRules модуля Settings (этап-1 готов); доски/правила/карточки — через публичный интерфейс модуля Kanban: порт IKanjStore (GetBoardAsync/AddCardAsync), статические чистые ColumnRules/BudgetNormalizer/AmountParser; ML — существующий порт IMlClient (Contracts); ИИ — новый порт IAiClassifier (Contracts/Integrations) с детерминированным LocalAiClassifier в Infrastructure (замена gRPC-клиентом ai-service на этапе 6). Адаптеры EF — в Deal.Infrastructure: PipelineStore, доработка KanbanStore (жёсткое удаление карточки чистит строки DedupEntries по LeadId), доработка LocalMlClient НЕ требуется (счётчики решений ml/ai инкрементирует сам модуль Pipeline в KV). HTTP — Deal.Api/Endpoints (MapPipelineEndpoints, /api/demo/ingest в MapDemoEndpoints); фоновые циклы — PipelineWorkerScheduler (2 с) и доработка StorageTickScheduler (чистка отсева). Публикации SSE — только из Api-слоя (Ruling 5 этапа 3): new_lead при создании карточки воркером, toast при автоочистке отсева; pipeline_stats НЕ публикуем (фронт его не слушает — Ruling 5/9).

Spec: docs/api/api-map.md §3.6 (L176186), §2 SSE (L3343), правила (L7–24; п.9 «экономия» L399, кривые места L390400, п.1 SSE L43); §4.5 очередь и отсев (L306–314), §4.1 карточка (L228257), §4.6 (L319341 — настройки обработки: stopPhrases/minLen/blockResumes/wantedType/budgetRequired*/autoArchive/ archiveAfterDays/aiEnabled/aiFilterEnabled/mlEnabled/domainKeywords/hireMarkers/resumeMarkers/levelTerms); docs/spec/ТЗ-дейл-новая-архитектура.md §5 «Обработка входящих» (L84–121), §7 «Вкладка „Обработка"» (L150161); roadmap (этап 4, L57–62); референс-семантика прототипа: backend/app/services/pipeline.py (целиком: enqueue L5385, stage1_plain L94124, clean_short/block L148193, _skip_no_budget L196218, compose_summary L225284, нормализация L294–341, контакты L344430, _store_lead L433514, локальные поля L661–798, воркер L8031183), backend/app/services/processing.py (целиком: record L66101, purge_expired L104117, clear_all/return_to_queue L120193, list_queue/list_rejected/stats L201320), backend/app/routers/processing_routes.py (целиком), backend/app/services/fts.py (целиком), backend/app/services/leads.py (L225247 _hard_delete/clear_col, L454504 tick_storage + notify, L509551 search), backend/app/services/ai.py (L261267 normalize_dedup, L316352 clean_budget/ budget_to_target), backend/app/services/ml_client.py (L2628 веса, L160162 is_enabled), backend/app/routers/dashboard_routes.py (L261284, L327337), backend/app/constants.py, backend/app/db.py (L8892 dedup, L226268 pipeline_msg/rejected_msgs); фронт: src/frontend/src/views/ProcessingView.vue (вся вкладка: счётчики L221–243, очередь L297400, отсев L402556, canReturn/return), src/frontend/src/store.js (pipeline-секция L12101343: loadPipelineQueue L12131223 limit=120, loadRejected L12261246 limit=80 offset, refreshPipelineStats L12491258, deleteRejectedItem L12841295, returnRejected L13001318, clearRejectedAll L13201333, SSE L676679 — обработчик pipeline_stats недостижим, api.js L62–104 слушает только 4 события), src/frontend/src/api.js, utils.js (tgSourceUrl); конвенции планов этапов 1–3 (файлы docs/superpowers/plans/2026-09-05-deal-stage{1,2,3}-*.md).

Global Constraints

  • Проект НЕ git; фиксация — отчёты задач task-N-report.md и progress.md в .superpowers/sdd/deal-stage4-pipeline/.
  • .NET 10 SDK, решение собирается 0 warnings / 0 errors (TreatWarningsAsErrors); dev-Postgres deal-postgres (:5433); curl-приёмка :5080 (scripts/build.sh/scripts/test.sh).
  • Код-стайл этапов 1–3: 1 тип = 1 файл; XML-doc на public-контракты; комментарии на русском; явные модификаторы; без регионов; без магических чисел (именованные константы); PascalCase-колонки БД; времена — DateTimeOffset (UTC) в БД, наружу epoch-ms; JSON camelCase; ошибки {"detail"}.
  • Модуль Pipeline — чистый: без EF и HTTP; зависимости — Deal.Modules.Settings (порт ISettingsStore, сервис IncomingRules), Deal.Modules.Kanban (порт IKanjStore, статические ColumnRules/BudgetNormalizer/ AmountParser, модели CardSnapshot/CardDto), Deal.Contracts (IMlClient, IAiClassifier). Реверс-зависимостей нет (Kanban/Settings о Pipeline не знают). Kanban НЕ получает ссылок на Pipeline — слияние статистик тика и жёсткое удаление dedup — в адаптерах/Api (Rulings 3/9).
  • LeadRadar-контейнеры, backend/, mlservice/, src/frontend/ НЕ трогаем. Vue-фронт не переписывается: формы JSON 1:1 с api-map; «кривые места» этапа 4: pipeline_stats у фронта недостижим (Ruling 9), GET /api/searchmessages: [].
  • Строки ошибок/тостов/причин — фиксированные из прототипа (см. задачи); новые строки — только для согласованных добавок (демо-ingest, Ruling 11).

Зафиксированные решения (Rulings этапа)

  • Ruling 1 (а) — миграция TenantPipeline и таблицы. Новая миграция TenantPipeline контекста TenantDbContext (папка I/Migrations/TenantDb, применяется провижинером ко всем схемам). Таблицы (PascalCase, владелец — модуль Pipeline; соответствие db.py L8892/L226268): QueueItems (= pipeline_msg), RejectedItems (= rejected_msgs), DedupEntries (= dedup). Колонки QueueItems: Id (p_, текст), DialogId, ChannelName/ChannelHandle/ChannelHue (дефолт #666), Text (≤6000), MsgId (long?, nullable), MsgAt, Status (new|filtered), Force (bool), CreatedAt, UpdatedAt; индекс (Status, CreatedAt). RejectedItems: Id (текст; детерминированный r_<dialog>_<msgId> при наличии dialog+msgId, иначе r_+hex — как processing.record L77; upsert ON CONFLICT (id) DO UPDATE), DialogId, MsgId (long?), ChannelName/ ChannelHandle/ChannelHue, Text (≤6000), Stage, Reason (≤500), Kw (≤200), Source, MsgAt, RejectedAt, Returned (bool), ReturnedAt (nullable), ReturnReason (≤500), SearchTsv (см. Ruling 6); индекс (RejectedAt)
    • GIN (SearchTsv). DedupEntries: Hash (текст, PK), LeadId (nullable, БЕЗ FK — «мягкая» ссылка на Cards, как прототип; чистка при жёстком удалении карточки — Ruling 3), CreatedAt. JSON-полей нет (все поля — плоские колонки); связи с Cards нет FK (журнал/отсев живут дольше карточки, конвенция Ruling 1 этапа 3). В той же миграции — FTS: Cards.SearchTsv и RejectedItems.SearchTsv (Ruling 6). Индексы/конфиги — 1 файл на сущность, эталон CardEntity+CardConfiguration.
  • Ruling 2 (б) — порт приёма сообщений и демо-ингвест. Приём — публичный scoped-сервис модуля PipelineIngestService.EnqueueAsync(QueuedMessage message, CancellationToken) (1:1 prototype enqueue L5385: trim текста, пустой текст/нет dialog → no-op; text[:6000]; msg_id-дубль-гвард на уровне адаптера SELECT 1 FROM QueueItems WHERE DialogId=? AND MsgId=? — защита от двойного события Telethon; id p_). На этапе 4 его вызывает ТОЛЬКО демо-эндпоинт POST /api/demo/ingest (Ruling 11; флаг DEAL_DEMO, иначе 404 «Демо-режим отключён»); этап 6 — gRPC-ингресс telegram-service вызовет тот же сервис (контракт стабилен, интерфейс не плодим — YAGNI). Разбор очереди — воркер (Ruling 8) + POST /api/admin/tick (Ruling 10), как прототип (pump L890–918 вызывается из _pipeline_loop и admin_tick L336).
  • Ruling 3 (в-1) — кто пишет карточку и доступ к Kanban. Карточку создаёт модуль Pipeline, но ТОЛЬКО через публичный интерфейс модуля-владельца Kanban (архитектура §5 L130–131): IKanjStore.AddCardAsync (CardSnapshot) + GetBoardAsync; проверка назначения колонки и «почему в колонке» — статические чистые ColumnRules.BoardAccepts/HasActiveRules/ComputeHits и BudgetNormalizer/AmountParser модуля Kanban (доступ к чистым помощникам владельца — не дублируем). Добавление ссылки Pipeline → Kanban цикла не создаёт (Kanban про Pipeline не знает). Подготовка полного CardSnapshotCardComposer в модуле Pipeline (перенос _store_lead L433514, Ruling 4). Жёсткое удаление карточки (Kanban DELETE /leads/{id}, clear-col, очистки тика) по контракту api-map §3.2 L92 — «leads+dedup+messages»: дорабатываем EF-адаптер KanbanStore (DeleteForeverAsync/PurgeAsync/ClearColAsync дополнительно удаляют DedupEntries WHERE LeadId=? — «сирота» не должна блокировать повторное создание, leads.py _hard_delete L229). Порт Kanban и его XML-doc обновляются (семантика «полное удаление»).
  • Ruling 4 (г) — карточка из сообщения (CardComposer). Перенос _store_lead (L433514) в чистый CardComposer модуля Pipeline: title = clean_short(raw.title, 140) или clean_short(text, 140); summary = compose_summary (блоки «О заявке» Компания→Формат→О задаче→Требования→Будет плюсом→Условия, 1:1 с cardPrompt и compose_summary L225284; локальный путь без структуры — «О задаче: …», _local_summary L294314; футер-хинты L288291) ≤2000 через clean_block; stack = normalize_stack ≤12 (L332341); бюджет: нормализованный из разбора (BudgetNormalizer.Normalize), иначе fallback из первой суммы AmountParser.Parse по исходнику/суммари (L459–468), конверсия один раз при поступлении — BudgetNormalizer.ToTarget (conversionOn/targetCurrency/ratesCache, USDT=USD, мок-фолбэк, как ConversionRecomputer/CardsService.LoadRatesAsync); контакты: ContactsQualifier.Build из разбора или текста (L389–421, ≤6, типы tg/phone/email/linkedin/whatsapp/site, отбрасывание ботов/сервисных t.me/ «постовых» сайтов L344386), primary_contact (tg→phone→whatsapp→email→linkedin→site, L424430, ≤200); ch-поля канала; sourceMsg = text[:4000]; sourceDialogId/sourceMsgId; prevCol=inbox; matchHits = ComputeHits доски, если назначена и прошла BoardAccepts (иначе колонка сбрасывается в inbox — страховка L449450); isVacancyKnown = признак ИИ. Создание: IKanjStore.AddCardAsync затем IPipelineStore.LinkDedupAsync(hash, cardId) (порядок как L512–513).
  • Ruling 5 (в-2) — ML/ИИ-ветки этапа 4. ML-слой вызывает существующий порт IMlClient.PredictAsync когда mlEnabled (не false) и не force (L966). Локальная модель не готова (LocalMlClient ready:false → predict {take:false,...}) — все сообщения уходят к ИИ-ветке; ветки «решил сам» реализуются ПОЛНОСТЬЮ 1:1 с L969–1061 (spam → отсев {source:ml, stage:spam_ml, reason:«ML уверен, что это спам/не заявка (score …)»}; доска → разрешена только не-suggested без активных правил, карточка в доску с локальными полями + типом ML + докладом terms в стек; typeDrop по wantedType; тип известен + aiEnabled false → карточка inbox) и покрываются юнит-тестами на fake-клиенте с ready:true (FakeMlClient в тестах расширяется). ИИ-слой: новый порт IAiClassifier (Contracts/Integrations; этап 6 заменит реализацию gRPC-клиентом ai-service) с record-DTO AiFilterResult {Pass, Reason, Skipped} и AiParsedLead (title/company/format/task/requirements/plus/conditions/stack/budget/contacts/is_vacancy/is_vacancy_known/ is_spam/board — структура классификации ТЗ §5 L104106 и ai.py classify). Этап 4 — детерминированный LocalAiClassifier (Infrastructure/Integrations): фильтр — всегда {pass:true, skipped:true} (реального ИИ-фильтра нет; при aiFilterEnabled=true это ветка «ИИ недоступен» прототипа L1103–1106; отсевы spam_ai/filter_ai недостижимы — их причины готовы для этапа 6); классификатор — локальный разбор ядра MessageParseCore (Ruling 4/Ruling 7: budget из AmountParser, контакты qualify, is_vacancy по hire-маркерам, is_vacancy_known=false, board=null — «смысловые колонки до ИИ не назначаем», L954958). aiEnabled=false → тот же локальный разбор напрямую (прототип L1081–1096), без вызова порта. Возврат (force): ИИ-фильтр пропускается (L1097–1100), вердикт «спам» ИИ отменяется (L1117–1121). Счётчики решений: KV mlDecisions/aiDecisions инкрементирует модуль Pipeline после pump (ml=mlStored+mlDrop, ai=aiStored+aiDrop, ml_client.track_decisions L153157) через ISettingsStore read-modify-write — LocalMlClient.StatusAsync их уже читает (этап 3), контракт IMlClient не меняется.
  • Ruling 6 (е) — FTS. Механизм — встроенный полнотекстовый поиск Postgres БЕЗ внешних расширений (pg_trgm и DuckDB-FTS НЕ нужны: LIKE-дополнение на объёмах этапа выполняется сканом, а русская морфология есть в конфигурации russian): в миграции TenantPipeline добавляются генерируемые колонки Cards.SearchTsv и RejectedItems.SearchTsv = to_tsvector('russian', coalesce(<текст.поля>,'')) (Cards: Title+Summary+SourceMsg+Contact — поля поиска leads L527529; Rejected: Text — fts.py _FTS_TARGETS L2327) STORED + GIN-индексы. Колонки авто-актуальны (аналог DuckDB «rebuild каждые сутки» не нужен). Поиск карточки /api/search?q= (q≥2) — один SQL: col != 'taken' AND (SearchTsv @@ plainto_tsquery('russian', q) OR lower(title/summary/source_msg/contact) LIKE '%q%'), порядок ts_rank DESC, ReceivedAt DESC, limit 12 — кандидаты FTS LIKE как в leads.search L509551 (messages:[] — api-map п.3). Поиск отсева GET /pipeline/rejected?q= — FTS-кандидаты (SearchTsv @@ plainto_tsquery) ∪ LIKE-дополнение по lower(text)/reason/kw/ch_name (processing.list_rejected L246277, лимиты limit*2 на каждую выборку, total = размер объединения, страницы по offset/limit ≤500). POST /admin/fts/rebuild — реальная идемпотентная обслуживающая операция FtsMaintenance.RebuildAsync: CREATE INDEX IF NOT EXISTS + ANALYZE обеих таблиц (самовосстановление индекса, если отсутствует), ответ {ok:true, ready:true}.
  • Ruling 7 (в-3) — чистое ядро разбора в модуле Pipeline. Перенос функций pipeline.py в чистые классы модуля MessageParseCore (1 тип = 1 файл): MessageTextCleaner (clean_short L148156 / clean_block L158193 — markdown-ссылки, **__~~, ||, голые URL, эмодзи-диапазоны, «C#»-защита, схлопывание, обрезка по границе), MessageListNormalizer(normalize_list L317329, normalize_stack L332341),ContactsQualifier(L344430: qualify_contact/build_contacts/primary_contact + регэкспы/наборы L345347, L597604),DedupHasher(normalize_dedup ai.py L261267:[^\wа-яё]+→ SHA1 hex),SummaryComposer(compose_summary + _local_summary + футер-хинты L288291),LocalFieldsParser(_local_fields L718798: меткиСтек/Грейд/Контакты/БюджетL591596 через_field_of-эквивалент, fallback-извлечения, заголовок, суть, is_vacancy по hireMarkers, is_vacancy_known=false, board=null) + словарь стоп-слов стека (_STOP_STACKL604–610), маркеры найма/грейда/резюме читаются из настроек (S) как в IncomingRules. Эти же классы используетLocalAiClassifier` (Infrastructure). Unit-тесты — на эталонных текстах (кейсы из devtests/e2e прототипа + примеры вакансий/заказов с контактами и бюджетами).
  • Ruling 8 (ж/з) — воркер, очистки, счётчики, SSE. PipelineWorkerService.PumpOnceAsync (модуль) — перенос _pump_unlocked L920–1183 (порядок строго 1:1): для status='new' (лимит 12): force? → stale-проверка (только не force; msgAt старше archiveAfterDays*суток при autoArchive=true → отсев {source:stale, stage:stale, reason:«сообщение старше N дн. (срок до автоархива) — не заводим в систему»}, строка удаляется, карточка НЕ создаётся — правка владельца «устаревшие не попадают в систему») → IncomingRules.CheckAsync (не прошёл → отсев {source:stop, stage=kind(length|stop|resume| type), reason, kw}, строка+dedup-claim удаляются) → дедуп (DedupHasher; хэш уже в DedupEntries → отсев {source:dup, stage:dup, reason:«сообщение уже в системе: карточка создана ранее или этот текст уже обрабатывается»}, удаление строки и её dedup-claim; иначе INSERT claim LeadId=null) → ML-слот (Ruling 5) → не решено → status='filtered'. Для status='filtered' (лимит 4): force? → stale (только не force) → aiEnabled=false: локальный разбор + no-budget(не force) → карточка inbox (счётчик aiStored — имя прототипа) ; aiEnabled=true: force → фильтр-пропуск; иначе IAiClassifier.FilterAsync (локально pass/skipped); ClassifyAsync; сбой/пустой разбор → локальный разбор (aiFail); is_spam → отсев {source:ai, stage:filter_ai|spam_ai, reason:«ИИ-фильтр: …»|«ИИ: не заявка — спам, реклама, скам или служебное сообщение»} + IMlClient.PushAsync(text, "spam", AI_WEIGHT 0.4); no-budget (не force) → отсев {source:stop, stage:budget, reason:«включён фильтр „не создавать карточку без суммы" — в тексте не указан бюджет»}; карточка (CardComposer) + new-лид-сигнал; обучение ML: ИИ назначил доску (не inbox и не-suggested, без активных правил) → PushAsync(text, boardId, 0.4); тип известен → PushAsync(text, "t:hire"|"t:order", 0.4) (L11551180). Результат pump — PipelinePumpResult: счётчики {staged, rulesStored, mlStored, mlDrop, typeDrop, aiStored, aiDrop, aiFail, noBudget} (1:1 имена wire-ключами admin/tick pipeline-словаря) + IReadOnlyList<CardDto> CreatedCards (для SSE, Ruling 9) + счётчики решений для KV. Воркер-гейт «не параллелить pump одного тенанта» — PipelinePumpGate (Api, singleton, Interlocked/ConcurrentDictionary; аналог asyncio.Lock L40). Очистка отсева: PipelineProcessingService.PurgeExpiredAsync (RejectedAt старше 3 суток, processing.purge_expired L104117) вызывается из тика (Ruling 10); полная ручная очистка — отдельный эндпоинт /rejected/clear. Счётчики вкладки — GET /pipeline/stats (queue.counts из QueueItems по status + rejected count). SSE этапа 4: pipeline_stats НЕ публикуем — api-map L43 фиксирует, что фронтовый openEvents() слушает только new_lead/toast/reminder_due/ system_status, а «Обработка» живёт на поллинге (ProcessingView reloadAll 2,6 с + store.js 60 с); публикуем: new_lead (полный CardDto; из Api после PumpOnce — admin/tick и PipelineWorkerScheduler, Ruling 5 этапа 3) и toast «Отсев очищен: N записей (3 дн.)» (trash) при ненулевой автоочистке (доработка StorageToastPublisher, notify_tick_stats L503504).
  • Ruling 9 (и) — интеграция с тиком/настройками без циклов. StorageTickService (Kanban) НЕ трогаем (purgedRejected=0 у него остаётся). Автоочистку отсева выполняет модуль Pipeline (PipelineProcessingService.PurgeExpiredAsync) в рамках тика: оркестрацию делает Api — POST /api/admin/tick вызывает Kanban-тик + purge-отсева + pump (Ruling 10), фоновый StorageTickScheduler — Kanban-тик + purge-отсева на каждый тенант; ответ тика объединяет статистику (storage = {…, purgedRejected} 1:1 с leads.tick_storage L488493). Публикация тостов — StorageToastPublisher. Настройки этапа-1 переиспользуют IncomingRules (Settings, scoped) и новые порции настроек читаются через ISettingsStore/SettingsKeys + дефолты SettingsDefaults (без дублирования каталога ключей). Спам-квоты/«системный отсев сверх stale|dup» в прототипе нет — НЕ реализуем (за этапом; roadmap §L57–62 трактуем как stale/dup source = «система», уже покрыто).
  • Ruling 10 — эндпоинты этапа и DI. Входят: 6 эндпоинтов /api/pipeline/* (api-map §3.6) — GET /stats, GET /queue (limit ≤500, дефолт 100; ответ {items, counts:{new,ai,total}, rejected}), GET /rejected (q/offset/limit ≤500; {items,total,offset,limit}), POST /rejected/clear → {ok:true, cleared}, DELETE /rejected/{rejId} → {ok:true}, POST /rejected/{rejId}/return {reason=""}{id, returned:true, returnedAt} (404 «Запись не найдена»; 400 «Сообщение уже возвращено в обработку»/«Повтор: карточка с таким текстом уже есть в системе — возвращать нечего»/«В записи нет текста сообщения»; при stage ∈ {spam_ml, spam_ai, filter_ai} — PushAsync(text,"spam",-1.0); строки очереди с force=true; запись отсева помечается returned/returnedAt/returnReason, НЕ удаляется) + демо-ingest POST /api/demo/ingest {text, dialogId?, channelName?, channelHandle?, channelHue?, msgId?, msgAt?}{ok:true, id, queue:{new,ai,total}} (400 «Текст сообщения пуст»; DEAL_DEMO guard). Возврат из отсева «мимо ML к ИИ» (ТЗ §5 L100) обеспечивает force. Модифицируются: POST /api/admin/tick (ответ 1:1 {storage, reminders:[], pipeline:<dict pump>, queue:int} + тосты + new_lead по созданным карточкам), POST /api/admin/fts/rebuild (Ruling 6). НЕ реализуем (фронт не вызывает, api-map п.9): admin/wipe|clear-cards|pump-gate, ml/learn|flush, /leads/{id}/seen; reclassify остаётся заглушкой Ruling 11 этапа 3. DI: AddPipelineModule() (модуль: Ingest/Processing/Worker/Rejects/ядра), адаптеры в AddDealPersistence (IPipelineStore → PipelineStore), AddDealIntegrations (+IAiClassifier → LocalAiClassifier); Program.cs — AddPipelineModule + MapPipelineEndpoints + hosted services (Ruling 8/10). Id-префиксы Pipeline — p_ (очередь), r_ (отсев; детерминированный вариант), хэш-ключ без префикса.
  • Ruling 11 (к) — демонстрация сквозного пути без telegram. Пресеты демо НЕ заводим: POST /api/demo/ingest принимает произвольный текст (детерминированная приёмка curl-текстами из Task 13: вакансия с бюджетом/контактами → карточка; короткое сообщение/стоп-фраза/резюме/чужой тип → отсев; одинаковый текст дважды → «повтор»; msgAt старше срока → «устарело»; без суммы при budgetRequiredHire=true → «нет суммы»). simulate-lead/age-lead этапа 3 не меняются. После ingest очередь разбирается фоном (2 с) или POST /api/admin/tick (детерминированно в curl). Вне этапа 4: реальные ai/telegram/ml-сервисы и gRPC (этап 6), Projects/reminder_due (этап 5), discovery, оператор/лимиты (этап 7), события pipeline_stats/boards_changed/leads_reclassified (недостижимы у фронта).

Задачи

Сокращения путей: P= src/core/Deal.Modules.Pipeline/, K= src/core/Deal.Modules.Kanban/, S= src/core/Deal.Modules.Settings/, C= src/core/Deal.Contracts/, I= src/core/Deal.Infrastructure/, A= src/core/Deal.Api/, T= src/core/tests/Deal.Tests.Unit/. Отчёты — task-N-report.md в .superpowers/sdd/deal-stage4-pipeline/.

Task 1: Миграция TenantPipeline — QueueItems/RejectedItems/DedupEntries + FTS-колонки

Files:

  • Create: I/Persistence/Entities/{QueueItemEntity,RejectedItemEntity,DedupEntryEntity}.cs и I/Persistence/{QueueItemConfiguration,RejectedItemConfiguration,DedupEntryConfiguration}.cs (поля/индексы Ruling 1; Text/Reason/Kw — text; времена — DateTimeOffset; SearchTsv — computed).
  • Modify: I/Persistence/Entities/CardEntity.cs + I/Persistence/CardConfiguration.cs — свойство SearchTsv (HasComputedColumnSql("to_tsvector('russian', coalesce(\"Title\",'')||' '||coalesce(\"Summary\",'')||' '||coalesce(\"SourceMsg\",'')||' '||coalesce(\"Contact\",''))", stored:true) + GIN-индекс) — Ruling 6.
  • Modify: I/Persistence/TenantDbContext.cs — DbSet'ы + ApplyConfiguration.
  • EF: миграция TenantPipeline для TenantDbContext (как TenantKanban: dotnet ef migrations add TenantPipeline --context TenantDbContext --output-dir Migrations/TenantDb --project src/core/Deal.Infrastructure --startup-project src/core/Deal.Api); старт Api применяет её к схеме дефолтного тенанта.

Источники: db.py L8892, L226268; Rulings 1/6; эталон: TenantKanban-миграция, CardEntity/Configuration.

Acceptance: build 0/0; dotnet test MarkerTests PASS; psql (search_path дефолтного тенанта): таблицы QueueItems/RejectedItems/DedupEntries + PK; Cards получила SearchTsv (generated, stored) и RejectedItems.SearchTsv; индексы IX_QueueItems_Status_CreatedAt, IX_RejectedItems_RejectedAt, GIN на SearchTsv (обоих таблиц); __TenantMigrationsHistory содержит TenantPipeline. Отчёт: task-1-report.md.

Task 2: Модуль Pipeline — DTO, словари отсева, порт IPipelineStore, реестр

Files:

  • Create: P/Application/Models/QueueItemDto.cs (§4.5 очередь L308: id/dialogId/msgId/text/status/ch{name, handle,hue}/msgAt/queuedAt), RejectedItemDto.cs (§4.5 отсев L310313: +stage/stageLabel/reason/kw/ source/sourceLabel/rejectedAt/returned/returnedAt/returnReason; наружу epoch-ms), QueueCountsDto.cs ({new,ai,total}), RejectRecord.cs (команда записи отсева: source/stage/reason/kw + канальные поля), QueuedMessage.cs (команда приёма: dialog/ch/msgId/text/msgAt/force), PipelinePumpResult.cs (счётчики Ruling 8 + CreatedCards), PipelineRejectConstants.cs (словари stage→stageLabel, source→sourceLabel, «система», Ruling 1/Ruling 9; processing.py L2646).
  • Create: P/Application/IPipelineStore.cs — порт (реализация — EF-адаптер Task 3): Queue (ExistsDuplicateAsync(dialogId,msgId), AddAsync, ListAsync(limit), CountByStatusAsync, SetStatusAsync, RemoveAsync); Rejects (UpsertAsync(RejectRecord) с детерминированным id, ListPageAsync(offset,limit), SearchIdsAsync(q, limitFts, limitLike) → упорядоченный список id, CountAsync, RemoveAsync, ClearAsync, PurgeExpiredAsync(olderThan), GetAsync(id), MarkReturnedAsync(id, reason, at)); Dedup (ExistsAsync(hash), ClaimAsync(hash), DeleteClaimAsync(hash) (только LeadId=null), LinkAsync(hash, cardId), DeleteByLeadAsync(cardId)).
  • Create: P/Application/PipelineIdPrefixes.cs (p_, r_) + переиспользование PrefixId (модуль Kanban) — при необходимости вынести общий генератор в SharedKernel (на усмотрение исполнителя, без дублирования).
  • Create: P/Application/PipelineModuleRegistrar.csAddPipelineModule() (регистрация сервисов задач 4/5/7/9 по мере появления). Modify: P/Deal.Modules.Pipeline.csproj — ProjectReference на Deal.Modules.Settings и Deal.Modules.Kanban.

Источники: api-map §4.5 L306314; processing.py L2646, L218320; Rulings 1/2/8/10.

Acceptance: build 0/0; DTO — record'ы (camelCase при сериализации); словари 1:1 (length→«короткое сообщение», …, dup→«повтор»; stop→«правила», ml→«ML», ai→«ИИ», stale|dup→«система»); MarkerTests PASS. Отчёт: task-2-report.md.

Task 3: EF-адаптер PipelineStore + DI + жёсткое удаление карточек (DedupEntries)

Files:

  • Create: I/Persistence/Repositories/PipelineStore.cs — реализация IPipelineStore на TenantDbContext (эталон KanbanStore.cs; AsNoTracking для чтений; маппинг вручную; времена ↔ epoch-ms наружу). Детали: ExistsDuplicateAsyncSELECT 1 FROM QueueItems WHERE DialogId=? AND MsgId=? (Ruling 2); добавление строки очереди — id p_ генерирует модуль; UpsertAsync для RejectedItems — raw SQL INSERT … ON CONFLICT (id) DO UPDATE SET … (processing.record L79101: детерминированный id r_<dialog>_<msgId> либо r_+hex; пустой текст — no-op); SearchIdsAsync — FTS-кандидаты plainto_tsquery('russian', q) по SearchTsv (rank DESC) + LIKE-дополнение по lower(text)/reason/kw/ch_name (limit*2 каждое), объединение без дублей (processing L252270); PurgeExpiredAsync/ClearAsync/RemoveAsync — по RejectedAt/безвозвратно; ClaimAsyncINSERT … ON CONFLICT DO NOTHING; DeleteClaimAsync удаляет только строки с LeadId IS NULL; LinkAsyncUPDATE DedupEntries SET LeadId=? WHERE Hash=?.
  • Modify: I/Persistence/Repositories/KanbanStore.cs — жёсткое удаление карточки (DeleteForeverAsync, PurgeAsync, ClearColAsync) дополнительно DELETE FROM DedupEntries WHERE LeadId=? (Ruling 3).
  • Modify: K/Application/IKanjStore.cs — XML-doc метода DeleteForeverAsync/PurgeAsync (семантика «Cards + комментарии + DedupEntries», Ruling 3).
  • Modify: I/ServiceCollectionExtensions.csAddScoped<IPipelineStore, PipelineStore>().

Источники: processing.py L66117, L246312; leads.py _hard_delete L225247; Rulings 1/3; эталон KanbanStore.cs/SettingsStore.cs.

Acceptance: build 0/0; unit (LocalMlClient-стиль не нужен — PipelineStore на EF покрывается curl/psql): upsert отсева дважды с тем же dialog+msgId → одна строка с обновлёнными полями; psql+curl — после Task 9. Отчёт: task-3-report.md.

Task 4: Чистое ядро разбора сообщения — cleaners, нормализация, контакты, dedup, «О заявке»

Files:

  • Create в P/Application/Parse/: MessageTextCleaner.cs (clean_short/clean_block L148193 + регэкспы/ наборы эмодзи/футер-хинты L131145, L288291), MessageListNormalizer.cs (normalize_list/normalize_stack L317341 + стоп-слова стека L604–610), ContactsQualifier.cs (L344430 + _contacts_from L666679, _norm_phone L661663), DedupHasher.cs (ai.py L261267), SummaryComposer.cs (compose_summary L225284 + _local_summary L294314), LocalFieldsParser.cs (_local_fields L718798 + _field_of L686698, метки L591–596, маркеры/токены L597612; hireMarkers/levelTerms/resumeMarkers — через ISettingsStore + SettingsDefaults, нормализация как в IncomingRules), AmountRangeBudgetFallback.cs (fallback бюджета из AmountParser.Parse, L459468).
  • Test: T/MessageParseCoreTests.cs — кейсы: markdown/URL/эмодзи-чистка, «C#» не режется, обрезка по границе; normalize_list «Java, Kotlin»/«;»-список; qualify: @user, @…bot → нет, t.me-ссылка, email, телефон +7, linkedin, site-спам (teletype.in → нет); build_contacts из текста (≤6, дедуп); dedup-хэш детерминирован (регистр/пунктуация не влияют, «Тест!» ≡ «тест»); compose_summary: блоки Компания→…→Условия в порядке; локальный путь «О задаче: …»; _local_fields на объявлении с метками «Стек:/Бюджет:/Контакты:» и без меток (fallback по тексту; is_vacancy по hire-маркерам, known=false, board=null).

Источники: pipeline.py L131341, L591798; ai.py L261267; Rulings 4/7.

Acceptance: dotnet test новых тестов PASS; build 0/0. Отчёт: task-4-report.md.

Task 5: PipelineService — приём (ingest), очередь, отсев, возврат, очистки, счётчики

Files:

  • Create: P/Application/PipelineIngestService.csEnqueueAsync(QueuedMessage) (Ruling 2: trim, no-op пустого текста/нет dialogId, text[:6000], msg_id-дубль-гвард, id p_, status=new, CreatedAt/UpdatedAt).
  • Create: P/Application/PipelineProcessingService.cs — запись отсева (Ruling 1/8: детерминированный upsert), чтение очереди (list_queue L218241: limit clamp 1..500), queue_counts (L207215), rejected_count, list_rejected (L246312: q-путь FTS+LIKE/страницы/лимиты, no-q путь по RejectedAt DESC), ReturnAsync (processing.return_to_queue L128193: 404 «Запись не найдена»; 400-строки Ruling 10; stage∈{spam_ml,spam_ai,filter_ai} → IMlClient.PushAsync(text,"spam",-1.0); пометка записи returned + return_reason; enqueue force=true с msg_at из записи), ClearAsync, DeleteAsync, PurgeExpiredAsync (3 суток от RejectedAt), stats (форма /pipeline/stats).
  • Test: T/PipelineProcessingServiceTests.cs + T/FakePipelineStore.cs (+ использование существующих FakeSettingsStore/FakeMlClient): ingest (trim/no-op/дубль-dialog+msgId), возврат: dup → 400-текст; повторный → 400; не найдена → 404-результат; спам-этап → PushAsync(spam, 1.0) вызван; очистки/счётчики.

Источники: pipeline.py L5385; processing.py L66193, L201320; processing_routes.py L1774; Rulings 2/8/10.

Acceptance: dotnet test PASS; build 0/0. Отчёт: task-5-report.md.

Task 6: Порт IAiClassifier + детерминированный LocalAiClassifier

Files:

  • Create: C/Integrations/IAiClassifier.cs + C/Integrations/Models/{AiFilterResultDto,AiParsedLeadDto, AiBudgetDto}.cs — порт Ruling 5: FilterAsync(string text, ct) и ClassifyAsync(string text, ct).
  • Create: I/Integrations/LocalAiClassifier.cs — реализация: фильтр всегда {pass:true, skipped:true} (aiFilterEnabled НЕ читает — выключатель обрабатывает воркер, как прототип filter_incoming L190192: выключен → skipped, включён при недоступном ИИ → pass+skipped, L11031106); классификатор — LocalFieldsParser (модуль Pipeline) → AiParsedLeadDto (title/summary/stack/budget из AmountParser/ BudgetNormalizer.Normalize/contacts через ContactsQualifier/is_vacancy/is_vacancy_known=false/board=null).
  • Modify: I/ServiceCollectionExtensions.csAddScoped<IAiClassifier, LocalAiClassifier>() (секция AddDealIntegrations).
  • Test: T/LocalAiClassifierTests.cs — фильтр-пропуск; классификатор детерминирован (одинаковый текст → одинаковый DTO); бюджет «до 2к$» → {from:null, to:2000, cur:USD}; контакты квалифицированы.

Источники: ai.py L188198, L316352; Rulings 5/7; эталон LocalColumnSuggester.cs (адаптер, зовущий модульное ядро).

Acceptance: dotnet test PASS; build 0/0. Отчёт: task-6-report.md.

Task 7: CardComposer — карточка из разобранного сообщения через публичный интерфейс Kanban

Files:

  • Create: P/Application/CardComposer.cs — сборка CardSnapshot из AiParsedLeadDto/локального разбора
    • метаданных сообщения (Ruling 4): title (clean 140), summary (compose_summary + clean_block 2000, fallback clean_short(text,2000)), stack ≤12, бюджет Normalize + fallback AmountParser по text/summary (первая сумма), ToTarget (conversionOn/targetCurrency/rates из ratesCache с мок-фолбэком, USDT=USD), contacts/primary contact, ch/source-поля, text[:4000], prevCol=inbox, isVacancy/Known. BuildAsync читает доску, если назначена (board): ColumnRules.BoardAccepts — иначе col=inbox; matchHits = ColumnRules.ComputeHits для прошедшей доски (иначе пусто).
  • Create: P/Application/PipelineCardWriter.cs — тонкая обёртка создания: PrefixId.New("l_")IKanjStore.AddCardAsync(snapshot)IPipelineStore.LinkDedupAsync(hash, cardId)store.GetCardAsync(cardId) (CardDto для SSE). (id l_ генерирует KanbanIdPrefixes — переиспользуем.)
  • Test: T/CardComposerTests.cs (FakeKanjStore/FakeSettingsStore): сборка полной карточки (блоки «О заявке», бюджет+conv, контакты, sourceMsg ≤4000); назначенная доска без правил → колонка доски + matchHits; доска с несовпадающими правилами → inbox (BoardAccepts-страховка); fallback-бюджет из текста; conv выключен (conversionOn=false) → conv-поля пусты.

Источники: pipeline.py L433514, L540586; rules.py board_accepts/hits (Kanban ColumnRules); Rulings 3/4.

Acceptance: dotnet test PASS; build 0/0. Отчёт: task-7-report.md.

Task 8: PipelineWorkerService — воркер pump (stale/правила/дедуп/ML/ИИ/карточка/обучение)

Files:

  • Create: P/Application/PipelineWorkerService.csPumpOnceAsync(newLimit=12, aiLimit=4) (Ruling 8, порядок 1:1 _pump_unlocked L9201183): проход new → проход filtered; создание карточек через PipelineCardWriter; отсевы через PipelineProcessingService; счётчики KV ml/ai инкремент после pump; возврат PipelinePumpResult (+CreatedCards). Зависимости: IPipelineStore, ISettingsStore, IncomingRules (Settings), IKanjStore, IMlClient, IAiClassifier, PipelineProcessingService, CardComposer, DedupHasher. Константы: PushWeightAi = 0.4, сроки из SettingsDefaults. Решения ML-ветки (Ruling 5) — на порту IMlClient: не готов/не уверен → filtered; spam/доска/тип — полные ветки.
  • Test: T/PipelineWorkerServiceTests.cs (+ доработка T/FakeMlClient.cs — настраиваемый ready/take/ label/type/terms; T/FakeAiClassifier.cs): (1) короткое → отсев length, строка удалена; (2) стоп-фраза → отсев stop с kw; (3) резюме → отсев resume; (4) dup: дважды один текст — второй отсев dup; (5) stale (msgAt старше срока, autoArchive=true) → отсев stale БЕЗ карточки; (6) вакансия с бюджетом → карточка inbox (aiStored=1, счётчики KV aiDecisions+1, CreatedCards=1); (7) no-budget при budgetRequiredHire → отсев budget, dedup-claim удалён; (8) ML ready+spam → отсев spam_ml + счётчик mlDecisions; (9) ML ready+доска (без правил, не suggested) → карточка в доску БЕЗ обучающего push (ML-путь не учит, L10171021); (10) ML ready+тип+aiEnabled=false → карточка inbox is_vacancy/known; (11) force: минует правила/ stale/no-budget и создаёт карточку; (12) ИИ-слот с fake-классификатором: доска назначена → BoardAccepts- страховка; is_spam → отсев spam_ai + Push(spam, 0.4); пустой разбор → локальный (aiFail).

Источники: pipeline.py L8031183; ml_client.py L2628; Rulings 5/8.

Acceptance: dotnet test новых тестов PASS; build 0/0. Отчёт: task-8-report.md.

Task 9: Эндпоинты /api/pipeline/* + /api/demo/ingest + DI + curl-приёмка

Files:

  • Create: A/Endpoints/PipelineEndpoints.cs (MapPipelineEndpoints): GET /api/pipeline/stats, GET /api/pipeline/queue?limit= (фронт шлёт 120; clamp 1..500; ответ {items, counts, rejected}), GET /api/pipeline/rejected?q=&offset=&limit= (clamp offset≥0/limit 1..500; {items,total,offset,limit}), POST /api/pipeline/rejected/clear{ok, cleared}, DELETE /api/pipeline/rejected/{rejId}{ok:true} (прототип delete_one L196198 всегда ok, 404 не шлём), POST /api/pipeline/rejected/{rejId}/return {reason} → 200 {id, returned:true, returnedAt} | 400 | 404 (детали Ruling 10). Статические сегменты до {rejId}; сессия 401 (эталон MlEndpoints/StorageEndpoints).
  • Create: A/Endpoints/RequestModels/ReturnReasonRequest.cs, PipelineIngestRequest.cs.
  • Modify: A/Endpoints/DemoEndpoints.csPOST /api/demo/ingest (флаг DEAL_DEMO; тело Ruling 11; 400 «Текст сообщения пуст»; вызов PipelineIngestService.EnqueueAsync; ответ {ok:true, id, queue:{new,ai,total}}).
  • Modify: A/Program.csAddPipelineModule(), MapPipelineEndpoints(); A/Deal.Api.csproj — ссылка на Deal.Modules.Pipeline.
  • Test: T/PipelineEndpointsContractsTests.cs — не нужен (endpoint-слои покрываются curl); достаточно существующих MarkerTests.

Контракт: api-map §3.6 L178186; §4.5; processing_routes.py.

Acceptance (curl admin/admin, DEAL_DEMO=1): stats/queue/rejected пустые формы; demo/ingest → очередь 1; ingest того же (dialogId+msgId) снова → очередь не растёт (гвард); queue?limit=120 — items/counts/rejected; rejected пуст; 401 без куки. Отчёт: task-9-report.md.

Task 10: POST /admin/tick и /admin/fts/rebuild реальные + SSE-тост отсева

Files:

  • Create: I/Services/FtsMaintenance.cs (или I/Persistence/Repositories/): RebuildAsync(context)CREATE INDEX IF NOT EXISTS для Cards(SearchTsv)/RejectedItems(SearchTsv) (raw SQL; имена — внутренние константы) + ANALYZE Cards/RejectedItems (Ruling 6).
  • Modify: A/Endpoints/StorageEndpoints.csAdminTickAsync: StorageTickService.TickAsync + PipelineProcessingService.PurgeExpiredAsync (merge в storage.purgedRejected) + PipelineWorkerService.PumpOnceAsync (один раз) + ответ {storage, reminders:[], pipeline:<PipelinePumpResult wire-dict>, queue:<count>} (dashboard_routes.py L327337); публикации: тосты StorageToastPublisher, new_lead на каждую карточку CreatedCards (Ruling 8/9). FtsRebuildAsync → FtsMaintenance + {ok:true, ready:true}.
  • Modify: A/Events/StorageToastPublisher.cs — ветка PurgedRejected > 0 → toast «Отсев очищен: N записей (3 дн.)» (trash) (notify_tick_stats L503504; тест T/StorageToastPublisherTests.cs дополняется).

Источники: dashboard_routes.py L261264, L327337; leads.py L486504; fts.py L4867; Rulings 6/8/10.

Acceptance: curl: demo/ingest вакансии → POST /api/admin/tick → pipeline содержит aiStored/созданную карточку (GET /leads), queue:0; после отсева (стоп-фраза) tick → pipeline-счётчики, /pipeline/rejected содержит запись; fts/rebuild → ok/ready. Отчёт: task-10-report.md.

Task 11: Фоновые циклы — PipelineWorkerScheduler (2 с) + purge-отсева в StorageTickScheduler

Files:

  • Create: A/PipelineWorkerScheduler.cs — IHostedService (эталон StorageTickScheduler/RatesRefreshScheduler): Timer 2 с; на каждое срабатывание — обход тенантов (системный репозиторий), на тенант — свой scope с ITenantContext; воркер-гейт A/PipelinePumpGate.cs (Interlocked per-tenant: admin/tick и цикл не разбирают очередь тенанта одновременно — аналог asyncio.Lock pipeline.py L40); после PumpOnce — публикация new_lead для CreatedCards (Ruling 8/9); try/catch + без подписчиков no-op.
  • Modify: A/Hosting/StorageTickScheduler.cs — после Kanban-тика каждого тенанта вызывать PipelineProcessingService.PurgeExpiredAsync и учесть в тостах (Ruling 8/9).
  • Modify: A/Program.csAddHostedService<PipelineWorkerScheduler>().

Источники: main.py _pipeline_loop/_storage_loop (L4353); pipeline.py L40; StorageTickScheduler.cs; Rulings 8/9/10.

Acceptance: build 0/0; запуск Api — лог без ошибок цикла; demo/ingest → в пределах ~5 с очередь разобрана (карточка в /leads или запись в /rejected) без ручного tick; psql: отсев со старым RejectedAt удаляется фоном (в пределах тика) + toast при подписанном SSE. Отчёт: task-11-report.md.

Task 12: Полнотекстовый поиск карточек — /api/search (FTS + LIKE)

Files:

  • Modify: K/Application/IKanjStore.cs + K/Application/Models/CardsQuery.cs (или новый метод): SearchCardsAsync(string q, int limit, CancellationToken) — упорядоченный список CardDto по Ruling 6.
  • Modify: I/Persistence/Repositories/KanbanStore.cs — реализация: raw SQL по Cards.SearchTsv (plainto_tsquery('russian') + ts_rank DESC, ReceivedAt DESC + LIKE по title/summary/source_msg/contact, col != 'taken', limit=12), затем полные CardDto (существующий маппинг/комментарии/time).
  • Modify: K/Application/CardsService.csSearchCardsAsync делегирует порту (старый перебор удаляется; поведение для q<2 — как сейчас, пусто).
  • Test: T/CardsServiceTests.cs — дополнить: вызов порта с q≥2/лимитом; q<2 → пусто (порт не зовётся).

Источники: leads.py search L509551; fts.py; api-map §3.2 L101; Rulings 6; этап 3 Task 7 (текущий LIKE-путь).

Acceptance: dotnet test PASS; build 0/0; curl (после Task 10-приёмки, карточки созданы): /api/search?q= <слово из title/source> → карточка; морфология «разработчик»/«разработчику» (по summary) → карточка (tsvector); messages: []. Отчёт: task-12-report.md.

Task 13: Финал этапа — интеграция и сквозная приёмка

  • scripts/build.sh/scripts/test.sh — успешны; dotnet build Deal.sln 0/0; все unit-тесты PASS (410 этапа 3 + новые).
  • Сквозной curl-сценарий (DEAL_DEMO=1, admin/admin): boot-группы → demo/ingest вакансии («Middle Python…, бюджет 16002200$, @crm_head, tg…», dialog demo_channel) → admin/tick → GET /leads: карточка l_… inbox (title/summary-«О заявке»/stack/budget/converted/contacts/ch/sourceMsg) → ingest короткого текста → tick → GET /pipeline/rejected: stageLabel «короткое сообщение», source «правила»; ingest текста со стоп-фразой (PATCH settings stopPhrases) → отсев stop c kw; повторный ingest того же текста вакансии → отсев dup (карточка уже есть); ingest без суммы при budgetRequiredHire=true → отсев «нет суммы»; ingest с msgAt старше archiveAfterDays → отсев «устарело» (карточки нет); GET /pipeline/queue?limit=120 — статусы new/filtered по ходу; GET /pipeline/stats — счётчики; GET /pipeline/rejected?q=<слово> (FTS) и ?q=<имя канала> (LIKE) → записи; DELETE /rejected/{id} → ok; POST /rejected/{id}/return {reason} (запись не dup/не returned) → возвращена в очередь (queue=1, запись returned=true), tick → карточка создана; повторный return той же записи → 400; POST /rejected/clear → {ok, cleared}; GET /api/search?q= по созданным карточкам; POST /admin/fts/rebuild → {ok,ready}; 401-проверки.
  • psql дефолтного тенанта: строки QueueItems/RejectedItems/DedupEntries; карточка ↔ dedup-связь (LeadId=карточка); удаление карточки (DELETE /leads/{id}) чистит DedupEntries; SearchTsv заполнены.
  • Обновить docs/technical/Техническая-документация-Дейл.md: раздел «Обработка/Pipeline» (таблицы этапа, эндпоинты /pipeline, демо-ingest, воркер-цикл 2 с, FTS, автоочистка отсева 3 дня, SSE-политика).
  • Отчёт task-13-report.md + финальная строка progress.md; roadmap-флаг «этап 4 выполнен».

Self-Review

  1. Spec coverage: ТЗ §5 (L84–121) — путь сообщения Tasks 5/8/10/11 (очередь→стоп-лист→дедуп→ML→ИИ→ карточка); этап-1 (длина/стоп-фразы/резюме/тип) — Task 8 через IncomingRules; «устаревшее» — Task 8; ML-слой (уверен — сам, иначе ИИ, возврат мимо ML) — Ruling 5/Task 8; ИИ-слой (фильтр/классификация/ колонка с проверкой правил) — Rulings 5/4, Tasks 6/7/8 (локальный детерминированный классификатор, реальный ИИ — этап 6); глобальные фильтры «без суммы» — Task 8; карточка (структура «О заявке», поля §5/§4.1, контакты-квалификация, конверсия) — Task 7; ТЗ §7 (L150161) — очередь/отсев/причины/ поиск/возврат/автоочистка/счётчики — Tasks 5/9 + Rulings 1/8; api-map §3.6/§4.5 — Tasks 2/3/5/9; §3.2 admin-tick/fts — Task 10; §2 SSE — Ruling 8/9; roadmap этап 4 — все задачи.
  2. Placeholder scan: заглушки — только согласованные: LocalMlClient (ready:false — ML-ветка «спит», ветки покрыты тестами на фейках), LocalAiClassifier (детерминированный до ai-service этапа 6; фильтр — pass+skipped, ветки отсева spam_ai/filter_ai готовы к этапу 6), demo-ingest (до telegram-этапа 6; контракт приёма — публичный сервис модуля), messages:[] в /api/search (api-map п.3), reclassify — заглушка этапа 3. Референсы на строки прототипа — точные; FIXME/TODO нет.
  3. Type consistency: Pipeline → Settings (порты/IncomingRules) и Pipeline → Kanban (IKanjStore + чистые помощники) — без циклов; Kanban не знает Pipeline; оркестрация тика и SSE — в Api (Ruling 5 этапа 3); FTS-колонки — в миграции TenantPipeline, владельцы таблиц не меняются (Kanban: Cards; Pipeline: QueueItems/RejectedItems/DedupEntries); контракт IMlClient не меняется (счётчики решений — KV через ISettingsStore); IAiClassifier в Contracts — подмена на gRPC этапа 6 без правки эндпоинтов; словари отсева/причины/тексты — 1:1 с прототипом; сущности/конфиги — конвенция TenantSettingEntity/CardEntity.
  4. Вне scope этапа 4: реальные ai/telegram/ml-сервисы и их gRPC-ингресс (этап 6; приём только demo- ingest), Projects/reminder_due (этап 5), discovery (этап 6), события pipeline_stats/boards_changed/ leads_reclassified (фронт не слушает — не публикуем), «спам-квоты»/новые глобальные exclude-настройки (в api-map/прототипе нет), admin/wipe|clear-cards|pump-gate, ml/learn|flush, /leads/{id}/seen, reclassify-реализация (этап 6), оператор/лимиты/аудит (этап 7).