Table of Contents
- Дейл (Deal) — Этап 4: Pipeline и «Обработка»: очередь, стоп-лист, дедуп, отсев, ML/ИИ-порты, FTS Implementation Plan
- Global Constraints
- Зафиксированные решения (Rulings этапа)
- Задачи
- Task 1: Миграция TenantPipeline — QueueItems/RejectedItems/DedupEntries + FTS-колонки
- Task 2: Модуль Pipeline — DTO, словари отсева, порт IPipelineStore, реестр
- Task 3: EF-адаптер PipelineStore + DI + жёсткое удаление карточек (DedupEntries)
- Task 4: Чистое ядро разбора сообщения — cleaners, нормализация, контакты, dedup, «О заявке»
- Task 5: PipelineService — приём (ingest), очередь, отсев, возврат, очистки, счётчики
- Task 6: Порт IAiClassifier + детерминированный LocalAiClassifier
- Task 7: CardComposer — карточка из разобранного сообщения через публичный интерфейс Kanban
- Task 8: PipelineWorkerService — воркер pump (stale/правила/дедуп/ML/ИИ/карточка/обучение)
- Task 9: Эндпоинты /api/pipeline/* + /api/demo/ingest + DI + curl-приёмка
- Task 10: POST /admin/tick и /admin/fts/rebuild реальные + SSE-тост отсева
- Task 11: Фоновые циклы — PipelineWorkerScheduler (2 с) + purge-отсева в StorageTickScheduler
- Task 12: Полнотекстовый поиск карточек — /api/search (FTS + LIKE)
- Task 13: Финал этапа — интеграция и сквозная приёмка
- Self-Review
Перенесено из репозитория (
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 (L176–186), §2 SSE (L33–43), правила (L7–24; п.9 «экономия» L399,
кривые места L390–400, п.1 SSE L43); §4.5 очередь и отсев (L306–314), §4.1 карточка (L228–257), §4.6
(L319–341 — настройки обработки: stopPhrases/minLen/blockResumes/wantedType/budgetRequired*/autoArchive/
archiveAfterDays/aiEnabled/aiFilterEnabled/mlEnabled/domainKeywords/hireMarkers/resumeMarkers/levelTerms);
docs/spec/ТЗ-дейл-новая-архитектура.md §5 «Обработка входящих» (L84–121), §7 «Вкладка „Обработка"»
(L150–161); roadmap (этап 4, L57–62); референс-семантика прототипа: backend/app/services/pipeline.py
(целиком: enqueue L53–85, stage1_plain L94–124, clean_short/block L148–193, _skip_no_budget L196–218,
compose_summary L225–284, нормализация L294–341, контакты L344–430, _store_lead L433–514, локальные
поля L661–798, воркер L803–1183), backend/app/services/processing.py (целиком: record L66–101,
purge_expired L104–117, clear_all/return_to_queue L120–193, list_queue/list_rejected/stats L201–320),
backend/app/routers/processing_routes.py (целиком), backend/app/services/fts.py (целиком),
backend/app/services/leads.py (L225–247 _hard_delete/clear_col, L454–504 tick_storage + notify,
L509–551 search), backend/app/services/ai.py (L261–267 normalize_dedup, L316–352 clean_budget/
budget_to_target), backend/app/services/ml_client.py (L26–28 веса, L160–162 is_enabled),
backend/app/routers/dashboard_routes.py (L261–284, L327–337), backend/app/constants.py,
backend/app/db.py (L88–92 dedup, L226–268 pipeline_msg/rejected_msgs);
фронт: src/frontend/src/views/ProcessingView.vue (вся вкладка: счётчики L221–243, очередь L297–400,
отсев L402–556, canReturn/return), src/frontend/src/store.js (pipeline-секция L1210–1343: loadPipelineQueue
L1213–1223 limit=120, loadRejected L1226–1246 limit=80 offset, refreshPipelineStats L1249–1258,
deleteRejectedItem L1284–1295, returnRejected L1300–1318, clearRejectedAll L1320–1333, SSE L676–679 —
обработчик 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-Postgresdeal-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/search→messages: []. - Строки ошибок/тостов/причин — фиксированные из прототипа (см. задачи); новые строки — только для согласованных добавок (демо-ingest, Ruling 11).
Зафиксированные решения (Rulings этапа)
- Ruling 1 (а) — миграция TenantPipeline и таблицы. Новая миграция
TenantPipelineконтекстаTenantDbContext(папкаI/Migrations/TenantDb, применяется провижинером ко всем схемам). Таблицы (PascalCase, владелец — модуль Pipeline; соответствие db.py L88–92/L226–268):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; upsertON 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.
- GIN
- Ruling 2 (б) — порт приёма сообщений и демо-ингвест. Приём — публичный scoped-сервис модуля
PipelineIngestService.EnqueueAsync(QueuedMessage message, CancellationToken)(1:1 prototype enqueue L53–85: trim текста, пустой текст/нет dialog → no-op; text[:6000]; msg_id-дубль-гвард на уровне адаптераSELECT 1 FROM QueueItems WHERE DialogId=? AND MsgId=?— защита от двойного события Telethon; idp_). На этапе 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 не знает). Подготовка полногоCardSnapshot—CardComposerв модуле Pipeline (перенос_store_leadL433–514, 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(L433–514) в чистыйCardComposerмодуля Pipeline: title = clean_short(raw.title, 140) или clean_short(text, 140); summary = compose_summary (блоки «О заявке» Компания→Формат→О задаче→Требования→Будет плюсом→Условия, 1:1 с cardPrompt и compose_summary L225–284; локальный путь без структуры — «О задаче: …», _local_summary L294–314; футер-хинты L288–291) ≤2000 через clean_block; stack = normalize_stack ≤12 (L332–341); бюджет: нормализованный из разбора (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/ «постовых» сайтов L344–386), primary_contact (tg→phone→whatsapp→email→linkedin→site, L424–430, ≤200); ch-поля канала; sourceMsg = text[:4000]; sourceDialogId/sourceMsgId; prevCol=inbox; matchHits = ComputeHits доски, если назначена и прошла BoardAccepts (иначе колонка сбрасывается в inbox — страховка L449–450); 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-DTOAiFilterResult {Pass, Reason, Skipped}иAiParsedLead(title/company/format/task/requirements/plus/conditions/stack/budget/contacts/is_vacancy/is_vacancy_known/ is_spam/board — структура классификации ТЗ §5 L104–106 и 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 — «смысловые колонки до ИИ не назначаем», L954–958). aiEnabled=false → тот же локальный разбор напрямую (прототип L1081–1096), без вызова порта. Возврат (force): ИИ-фильтр пропускается (L1097–1100), вердикт «спам» ИИ отменяется (L1117–1121). Счётчики решений: KVmlDecisions/aiDecisionsинкрементирует модуль Pipeline после pump (ml=mlStored+mlDrop, ai=aiStored+aiDrop, ml_client.track_decisions L153–157) через 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 L527–529; Rejected: Text — fts.py_FTS_TARGETSL23–27)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 L509–551 (messages:[]— api-map п.3). Поиск отсеваGET /pipeline/rejected?q=— FTS-кандидаты (SearchTsv @@ plainto_tsquery) ∪ LIKE-дополнение поlower(text)/reason/kw/ch_name(processing.list_rejected L246–277, лимиты 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 L148–156 / clean_block L158–193 — markdown-ссылки, **__~~, ||, голые URL, эмодзи-диапазоны, «C#»-защита, схлопывание, обрезка по границе),MessageListNormalizer(normalize_list L317–329, normalize_stack L332–341),ContactsQualifier(L344–430: qualify_contact/build_contacts/primary_contact + регэкспы/наборы L345–347, L597–604),DedupHasher(normalize_dedup ai.py L261–267:[^\wа-яё]+→ SHA1 hex),SummaryComposer(compose_summary + _local_summary + футер-хинты L288–291),LocalFieldsParser(_local_fields L718–798: меткиСтек/Грейд/Контакты/БюджетL591–596 через_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_unlockedL920–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)(L1155–1180). Результат 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 L104–117) вызывается из тика (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 L503–504). - 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 L488–493). Публикация тостов — 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, НЕ удаляется) + демо-ingestPOST /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 L88–92, L226–268; 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 отсев L310–313: +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 L26–46). - 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.cs—AddPipelineModule()(регистрация сервисов задач 4/5/7/9 по мере появления). Modify:P/Deal.Modules.Pipeline.csproj— ProjectReference наDeal.Modules.SettingsиDeal.Modules.Kanban.
Источники: api-map §4.5 L306–314; processing.py L26–46, L218–320; 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 наружу). Детали:ExistsDuplicateAsync—SELECT 1 FROM QueueItems WHERE DialogId=? AND MsgId=?(Ruling 2); добавление строки очереди — idp_генерирует модуль;UpsertAsyncдля RejectedItems — raw SQLINSERT … ON CONFLICT (id) DO UPDATE SET …(processing.record L79–101: детерминированный idr_<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 L252–270);PurgeExpiredAsync/ClearAsync/RemoveAsync— по RejectedAt/безвозвратно;ClaimAsync—INSERT … ON CONFLICT DO NOTHING;DeleteClaimAsyncудаляет только строки сLeadId IS NULL;LinkAsync—UPDATE 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.cs—AddScoped<IPipelineStore, PipelineStore>().
Источники: processing.py L66–117, L246–312; leads.py _hard_delete L225–247; 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 L148–193 + регэкспы/ наборы эмодзи/футер-хинты L131–145, L288–291),MessageListNormalizer.cs(normalize_list/normalize_stack L317–341 + стоп-слова стека L604–610),ContactsQualifier.cs(L344–430 + _contacts_from L666–679, _norm_phone L661–663),DedupHasher.cs(ai.py L261–267),SummaryComposer.cs(compose_summary L225–284 + _local_summary L294–314),LocalFieldsParser.cs(_local_fields L718–798 + _field_of L686–698, метки L591–596, маркеры/токены L597–612; hireMarkers/levelTerms/resumeMarkers — через ISettingsStore + SettingsDefaults, нормализация как в IncomingRules),AmountRangeBudgetFallback.cs(fallback бюджета изAmountParser.Parse, L459–468). - 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 L131–341, L591–798; ai.py L261–267; Rulings 4/7.
Acceptance: dotnet test новых тестов PASS; build 0/0. Отчёт: task-4-report.md.
Task 5: PipelineService — приём (ingest), очередь, отсев, возврат, очистки, счётчики
Files:
- Create:
P/Application/PipelineIngestService.cs—EnqueueAsync(QueuedMessage)(Ruling 2: trim, no-op пустого текста/нет dialogId, text[:6000], msg_id-дубль-гвард, idp_, status=new, CreatedAt/UpdatedAt). - Create:
P/Application/PipelineProcessingService.cs— запись отсева (Ruling 1/8: детерминированный upsert), чтение очереди (list_queue L218–241: limit clamp 1..500), queue_counts (L207–215), rejected_count, list_rejected (L246–312: q-путь FTS+LIKE/страницы/лимиты, no-q путь по RejectedAt DESC),ReturnAsync(processing.return_to_queue L128–193: 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 L53–85; processing.py L66–193, L201–320; processing_routes.py L17–74; 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 L190–192: выключен → skipped, включён при недоступном ИИ → pass+skipped, L1103–1106); классификатор —LocalFieldsParser(модуль Pipeline) →AiParsedLeadDto(title/summary/stack/budget из AmountParser/ BudgetNormalizer.Normalize/contacts через ContactsQualifier/is_vacancy/is_vacancy_known=false/board=null). - Modify:
I/ServiceCollectionExtensions.cs—AddScoped<IAiClassifier, LocalAiClassifier>()(секция AddDealIntegrations). - Test:
T/LocalAiClassifierTests.cs— фильтр-пропуск; классификатор детерминирован (одинаковый текст → одинаковый DTO); бюджет «до 2к$» → {from:null, to:2000, cur:USD}; контакты квалифицированы.
Источники: ai.py L188–198, L316–352; 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для прошедшей доски (иначе пусто).
- метаданных сообщения (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.
- Create:
P/Application/PipelineCardWriter.cs— тонкая обёртка создания:PrefixId.New("l_")→IKanjStore.AddCardAsync(snapshot)→IPipelineStore.LinkDedupAsync(hash, cardId)→store.GetCardAsync(cardId)(CardDto для SSE). (idl_генерирует KanbanIdPrefixes — переиспользуем.) - Test:
T/CardComposerTests.cs(FakeKanjStore/FakeSettingsStore): сборка полной карточки (блоки «О заявке», бюджет+conv, контакты, sourceMsg ≤4000); назначенная доска без правил → колонка доски + matchHits; доска с несовпадающими правилами → inbox (BoardAccepts-страховка); fallback-бюджет из текста; conv выключен (conversionOn=false) → conv-поля пусты.
Источники: pipeline.py L433–514, L540–586; 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.cs—PumpOnceAsync(newLimit=12, aiLimit=4)(Ruling 8, порядок 1:1_pump_unlockedL920–1183): проход 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-путь не учит, L1017–1021); (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 L803–1183; ml_client.py L26–28; 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 L196–198 всегда 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.cs—POST /api/demo/ingest(флаг DEAL_DEMO; тело Ruling 11; 400 «Текст сообщения пуст»; вызовPipelineIngestService.EnqueueAsync; ответ{ok:true, id, queue:{new,ai,total}}). - Modify:
A/Program.cs—AddPipelineModule(),MapPipelineEndpoints();A/Deal.Api.csproj— ссылка наDeal.Modules.Pipeline. - Test:
T/PipelineEndpointsContractsTests.cs— не нужен (endpoint-слои покрываются curl); достаточно существующих MarkerTests.
Контракт: api-map §3.6 L178–186; §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.cs—AdminTickAsync:StorageTickService.TickAsync+PipelineProcessingService.PurgeExpiredAsync(merge вstorage.purgedRejected) +PipelineWorkerService.PumpOnceAsync(один раз) + ответ{storage, reminders:[], pipeline:<PipelinePumpResult wire-dict>, queue:<count>}(dashboard_routes.py L327–337); публикации: тосты 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 L503–504; тестT/StorageToastPublisherTests.csдополняется).
Источники: dashboard_routes.py L261–264, L327–337; leads.py L486–504; fts.py L48–67; 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.cs—AddHostedService<PipelineWorkerScheduler>().
Источники: main.py _pipeline_loop/_storage_loop (L43–53); 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.cs—SearchCardsAsyncделегирует порту (старый перебор удаляется; поведение для q<2 — как сейчас, пусто). - Test:
T/CardsServiceTests.cs— дополнить: вызов порта с q≥2/лимитом; q<2 → пусто (порт не зовётся).
Источники: leads.py search L509–551; 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.sln0/0; все unit-тесты PASS (410 этапа 3 + новые).- Сквозной curl-сценарий (DEAL_DEMO=1, admin/admin): boot-группы → demo/ingest вакансии («Middle Python…, бюджет 1600–2200$, @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
- 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 (L150–161) — очередь/отсев/причины/ поиск/возврат/автоочистка/счётчики — 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 — все задачи.
- 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 нет. - 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.
- Вне 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).