using System.Globalization; using Deal.Contracts.Integrations.Abstractions; using Deal.Contracts.Integrations.Models; using Deal.Modules.Cards.Application.Models; using Deal.Modules.Cards.Application.Sources; using Deal.Modules.Kanban.Application.Abstractions; using Deal.Modules.Kanban.Application.Models; using Deal.Modules.Kanban.Application.Services; using Deal.Modules.Pipeline.Application.Abstractions; using Deal.Modules.Pipeline.Application.Models; using Deal.SharedKernel.Errors; namespace Deal.Modules.Pipeline.Application.Services; /// /// Ручная проверка/разметка ML на записях источников. /// public sealed class MlReviewService( IPipelineStore pipelineStore, ICardStore cardStore, CardsService cards, PipelineProcessingService processing, IMlClient mlClient) { /// /// Минимум сообщений в выборке кандидатов. /// public const int MinCandidates = 1; /// /// Максимум сообщений в выборке кандидатов. /// public const int MaxCandidates = 60; // Размер одного чтения из очереди/отсева при объединении кандидатов. private const int MaxScan = 500; // Имя сущности для текста ошибки «не найдено». private const string MessageEntityName = "Исходное сообщение"; private const int TextPreviewLength = 600; /// /// Вердикт кандидата /// public const string VerdictCard = "card"; /// /// Вердикт кандидата /// public const string VerdictRejected = "rejected"; /// /// Вердикт кандидата /// public const string VerdictQueued = "queued"; /// /// Действие: пропустить без обучения. /// public const string ActionSkip = "skip"; /// /// Действие: спам — учим ML и /// public const string ActionSpam = "spam"; /// /// Префикс действия «в колонку» /// public const string ActionBoardPrefix = "board:"; /// /// 400 apply: неизвестная доска-цель. /// public const string UnknownBoardDetail = "Неизвестная доска"; /// /// 400 apply: неизвестное действие. /// public const string UnknownActionDetail = "Неизвестное действие"; // Причина отсева при ручной разметке «спам» ещё не обработанного сообщения. private const string ManualSpamReason = "ручная разметка ML: спам"; // Этап отсева при ручной разметке «спам» (отсев решением ML). private const string ManualSpamStage = "spam_ml"; // Источник решения при ручной разметке. private const string ManualSource = "ml"; private const double UserPushWeight = 1.0; /// /// Отбирает записи-кандидаты для проверки ML по источнику и/или размеру выборки. /// /// Оригинал источника (OriginRef); пусто — выборка по всем источникам тенанта. /// Сколько последних записей вернуть (кламп 1.., дефолт вызывающего). /// Кандидаты (свежие первыми): текст, текущий вердикт и мнение ML по каждому. public async Task> CandidatesAsync( string? dialogId, int limit, CancellationToken ct) { int take = Math.Clamp(limit, MinCandidates, MaxCandidates); string dialog = (dialogId ?? string.Empty).Trim(); IReadOnlyList queue = await pipelineStore.ListAsync(status: null, MaxScan, ct); IReadOnlyList rejected = await pipelineStore.ListPageAsync(offset: 0, MaxScan, ct); IReadOnlyList cardList = await cardStore.ListCardsAsync(new CardsQuery(null), ct); // Объединение по ключу дедупа источника: очередь → отсев → карточка (последняя перекрывает предыдущие). var merged = new Dictionary(StringComparer.Ordinal); foreach (QueueItemDto row in queue) { if (!IsMessage(row.Source) || !MatchesDialog(dialog, row.Source.OriginRef)) { continue; } merged[row.Source.DedupeKey()] = BuildQueued(row); } foreach (RejectedItemDto row in rejected) { if (!IsMessage(row.Source) || !MatchesDialog(dialog, row.Source.OriginRef)) { continue; } merged[row.Source.DedupeKey()] = BuildRejected(row); } foreach (CardDto card in cardList) { if (!IsMessage(card.Source) || !MatchesDialog(dialog, card.Source.OriginRef)) { continue; } merged[card.Source.DedupeKey()] = BuildCard(card); } List ordered = merged.Values .OrderByDescending(candidate => candidate.Time ?? 0) .Take(take) .ToList(); var withPredictions = new List(ordered.Count); foreach (MlCandidateDto candidate in ordered) { withPredictions.Add(candidate with { Pred = await PredictSafelyAsync(candidate.Text, ct) }); } return withPredictions; } /// /// Применяет ручное решение по записи источника /// /// Оригинал источника (OriginRef) записи. /// Внешний id записи в источнике. /// Действие: skip | spam | board:<id>. /// Результат решения. /// Исходная запись не найдена. public async Task ApplyAsync( string dialogId, long msgId, string? action, CancellationToken ct) { string normalized = (action ?? string.Empty).Trim(); string dialog = (dialogId ?? string.Empty).Trim(); string externalId = msgId.ToString(CultureInfo.InvariantCulture); SourceRef source = await ResolveSourceAsync(dialog, externalId, ct) ?? throw new NotFoundException(MessageEntityName, externalId); CardDto? card = await cardStore.GetCardBySourceAsync(source, ct); string? text = await FindTextAsync(dialog, externalId, card, ct); if (string.IsNullOrWhiteSpace(text)) { throw new NotFoundException(MessageEntityName, externalId); } if (normalized == ActionSkip) { return new MlApplyResult(Error: null, Ok: true, Learned: false, Moved: null, LeadId: null); } if (normalized == ActionSpam) { return await ApplySpamAsync(dialog, externalId, card, text, ct); } if (normalized.StartsWith(ActionBoardPrefix, StringComparison.Ordinal)) { string boardId = normalized[ActionBoardPrefix.Length..].Trim(); return await ApplyBoardAsync(boardId, card, text, ct); } return new MlApplyResult(UnknownActionDetail, Ok: false, Learned: false, Moved: null, LeadId: null); } // Действие «спам»: карточку — в корзину (с обучением), запись из очереди — в отсев; иначе учим ML. // dialog: Оригинал источника. // externalId: Внешний id записи. // card: Карточка сообщения (null — сообщение не становилось карточкой). // text: Текст сообщения. // ct: Токен отмены. // Возвращает: Результат решения. private async Task ApplySpamAsync( string dialog, string externalId, CardDto? card, string text, CancellationToken ct) { if (card is not null) { // TrashCardAsync(teach:true) сам шлёт обучающий сигнал «спам» — второй сигнал не нужен. CardDto? trashed = await cards.TrashCardAsync(card.Id, teach: true, ct); return new MlApplyResult(null, Ok: true, Learned: true, Moved: "trash", LeadId: trashed?.Id ?? card.Id); } await mlClient.PushAsync(text, MlLearningLabels.Spam, UserPushWeight, ct); // Сообщение ещё в очереди — отсеиваем его (решение пользователя), снимая строку. QueueItemDto? queued = await FindQueuedAsync(dialog, externalId, ct); if (queued is not null) { await processing.RejectAsync(new RejectRecord { Source = queued.Source, Content = queued.Content, Text = queued.Text, MsgAtMs = queued.MsgAtMs, DecidedBy = ManualSource, Stage = ManualSpamStage, Reason = ManualSpamReason, Kw = string.Empty, }, ct); await pipelineStore.RemoveAsync(queued.Id, ct); } return new MlApplyResult(null, Ok: true, Learned: true, Moved: null, LeadId: null); } // Действие «в колонку»: карточку — переносим, уже в колонке — только учим; иначе учим ML. // boardId: Id колонки-цели (inbox или b_...). // card: Карточка сообщения (null — сообщение не становилось карточкой). // text: Текст сообщения. // ct: Токен отмены. // Возвращает: Результат решения. private async Task ApplyBoardAsync( string boardId, CardDto? card, string text, CancellationToken ct) { if (boardId != CardIds.Inbox && await cardStore.GetContainerAsync(boardId, ct) is null) { return new MlApplyResult(UnknownBoardDetail, Ok: false, Learned: false, Moved: null, LeadId: null); } if (card is null) { await mlClient.PushAsync(text, boardId, UserPushWeight, ct); return new MlApplyResult(null, Ok: true, Learned: true, Moved: null, LeadId: null); } if (card.Col == boardId) { await mlClient.PushAsync(text, boardId, UserPushWeight, ct); return new MlApplyResult(null, Ok: true, Learned: true, Moved: null, LeadId: card.Id); } // MoveDashboardCardAsync сам учит колонку (toCol ≠ inbox) — второй сигнал не нужен. CardResultDto moved = await cards.MoveDashboardCardAsync(card.Id, boardId, ct); if (moved.Error is not null) { return new MlApplyResult(moved.Error, Ok: false, Learned: false, Moved: null, LeadId: card.Id); } return new MlApplyResult(null, Ok: true, Learned: true, Moved: boardId, LeadId: card.Id); } // Ссылка на источник записи: очередь → карточка → отсев. // dialog: Оригинал источника. // externalId: Внешний id записи. // ct: Токен отмены. // Возвращает: Ссылку на источник либо null, если записи нет. private async Task ResolveSourceAsync( string dialog, string externalId, CancellationToken ct) { QueueItemDto? queued = await FindQueuedAsync(dialog, externalId, ct); if (queued is not null) { return queued.Source; } IReadOnlyList cardList = await cardStore.ListCardsAsync(new CardsQuery(null), ct); foreach (CardDto card in cardList) { if (MatchesSource(card.Source, dialog, externalId)) { return card.Source; } } IReadOnlyList rejected = await pipelineStore.ListPageAsync(offset: 0, MaxScan, ct); foreach (RejectedItemDto row in rejected) { if (MatchesSource(row.Source, dialog, externalId)) { return row.Source; } } return null; } // Текст исходного сообщения: текст карточки, иначе текст строки очереди/записи отсева. // dialog: Оригинал источника. // externalId: Внешний id записи. // card: Карточка сообщения (уже прочитана вызывающим). // ct: Токен отмены. // Возвращает: Непустой текст либо null, если записи нет ни в одном источнике. private async Task FindTextAsync( string dialog, string externalId, CardDto? card, CancellationToken ct) { if (card is not null && !string.IsNullOrWhiteSpace(card.Content.Text)) { return card.Content.Text; } QueueItemDto? queued = await FindQueuedAsync(dialog, externalId, ct); if (queued is not null && !string.IsNullOrWhiteSpace(queued.Text)) { return queued.Text; } IReadOnlyList rejected = await pipelineStore.ListPageAsync(offset: 0, MaxScan, ct); foreach (RejectedItemDto row in rejected) { if (MatchesSource(row.Source, dialog, externalId)) { return row.Text; } } return null; } // Строка очереди записи (для отсева при ручной разметке «спам»). // dialog: Оригинал источника. // externalId: Внешний id записи. // ct: Токен отмены. // Возвращает: Строка очереди либо null — сообщение уже обработано/не в очереди. private async Task FindQueuedAsync( string dialog, string externalId, CancellationToken ct) { IReadOnlyList queue = await pipelineStore.ListAsync(status: null, MaxScan, ct); foreach (QueueItemDto row in queue) { if (MatchesSource(row.Source, dialog, externalId)) { return row; } } return null; } // Прогноз ML по тексту с защитой от сбоя (недоступный сервис — кандидат без мнения). // text: Текст сообщения. // ct: Токен отмены. // Возвращает: Мнение ML либо null при сбое. private async Task PredictSafelyAsync(string text, CancellationToken ct) { if (string.IsNullOrWhiteSpace(text)) { return null; } try { MlPredictResultDto result = await mlClient.PredictAsync(text, ct); return new MlCandidatePredictionDto(result.Take, result.Label, result.Scores); } catch (Exception) { return null; } } // Идентифицирована ли запись источника (есть внешний id). // source: Ссылка на источник. // Возвращает: True — запись адресуема как сообщение источника. private static bool IsMessage(SourceRef source) => source.ExternalId is { Length: > 0 }; // Соответствует ли источник фильтру оригинала (пустой фильтр — все источники). // filter: Запрошенный оригинал (пусто — без фильтра). // originRef: Оригинал источника записи. // Возвращает: True — кандидат подходит. private static bool MatchesDialog(string filter, string? originRef) => filter.Length == 0 || string.Equals(filter, originRef ?? string.Empty, StringComparison.Ordinal); // Совпадает ли источник с оригиналом и внешним id записи. // source: Ссылка на источник. // dialog: Оригинал источника. // externalId: Внешний id записи. // Возвращает: True — источник указывает на ту же запись. private static bool MatchesSource(SourceRef source, string dialog, string externalId) => string.Equals(source.OriginRef ?? string.Empty, dialog, StringComparison.Ordinal) && string.Equals(source.ExternalId ?? string.Empty, externalId, StringComparison.Ordinal); // Кандидат из строки очереди (вердикт queued). // row: Строка очереди. // Возвращает: Кандидат. private static MlCandidateDto BuildQueued(QueueItemDto row) => new() { Source = row.Source, Content = row.Content, Text = Truncate(row.Text), Time = row.MsgAtMs == 0 ? null : row.MsgAtMs, Lead = false, Verdict = VerdictQueued, Stage = row.Status, }; // Кандидат из записи отсева (вердикт rejected). // row: Запись отсева. // Возвращает: Кандидат. private static MlCandidateDto BuildRejected(RejectedItemDto row) => new() { Source = row.Source, Content = row.Content, Text = Truncate(row.Text), Time = row.MsgAtMs == 0 ? null : row.MsgAtMs, Lead = false, Verdict = VerdictRejected, Stage = row.Stage, Reason = row.Reason, }; // Кандидат из карточки (вердикт card). // card: Карточка. // Возвращает: Кандидат. private static MlCandidateDto BuildCard(CardDto card) => new() { Source = card.Source, Content = card.Content, Text = Truncate(card.Content.Text ?? string.Empty), Time = card.ReceivedAtMs == 0 ? null : card.ReceivedAtMs, Lead = true, Verdict = VerdictCard, Col = card.Col, }; // Обрезает текст кандидата до TextPreviewLength символов. // text: Исходный текст. // Возвращает: Обрезанный текст. private static string Truncate(string text) => text.Length <= TextPreviewLength ? text : text[..TextPreviewLength]; }