using Deal.Api.Events;
using Deal.Modules.Kanban.Application;
using Deal.Modules.Kanban.Application.Models;
using Deal.Modules.Pipeline.Application;
using Deal.Modules.Pipeline.Application.Models;
namespace Deal.Api;
///
/// Оркестратор ручного тика POST /api/admin/tick (план Tasks 10–11, Ruling 8/9; dashboard_routes.py admin_tick L327–337).
///
///
/// Api-слой объединяет сервисы модулей (Kanban тик правил хранения + Pipeline очистка отсева и pump +
/// Projects проверка напоминаний «Отложено») и публикует SSE (Ruling 5/8/9 — публикации только из Api;
/// модули остаются чистыми). Порядок 1:1 с прототипом:
/// (1) — автоархив и очистки архива/корзины;
/// (2) — отсев старше 3 суток (tick_storage L485–493),
/// результат вливается в storage.purgedRejected (Ruling 9);
/// (3) SSE-тосты статистики (, notify_tick_stats L496–504) — до pump, как в
/// прототипе (L333);
/// (4) проверка наступивших напоминаний (admin_tick L334,
/// check_reminders L264–282; план Task 11, Ruling 3): «выстрелившие» {id,title,containerId} помечены fired и
/// публикуются SSE reminder_due (Ruling 8 — toast НЕ шлём, у фронта модалка ReminderNotice); сбой
/// проверки НЕ роняет тик: лог + reminders ответа пуст;
/// (5) под общим воркер-гейтом тенанта (Task 10/11): pump
/// одного тенанта выполняет либо ручной тик, либо фоновый цикл — при занятом гейте проход пропускается;
/// сбой pump НЕ роняет тик: исключение логируется, pipeline ответа пуст ({} как при занятом локе прототипа
/// L901–902), очередь остаётся до следующего тика/фонового цикла. Операция отмены (OCE) пробрасывается — запрос прерван;
/// (6) SSE new_card по каждой созданной карточке (Ruling 8/9; полный CardDto, как publish из Api);
/// (7) queue = строк очереди после pump (queue_len L337). Ответ — .
///
/// Тик правил хранения канбана (StorageTickService модуля Kanban).
/// Очистка отсева и счётчики очереди (модуль Pipeline).
/// Один проход pump по очереди входящих (модуль Pipeline).
/// Проверка наступивших напоминаний «Отложено» (CardsService, Ruling 3).
/// Публикатор SSE-тостов статистики тика (общий с фоновым циклом Task 11).
/// SSE-брокер канала тенанта (публикация reminder_due/new_card).
/// Общий воркер-гейт pump тенанта (singleton; общий с фоновым циклом Task 11).
/// Логгер сбоя проверки напоминаний/pump (тик продолжается без этих веток).
public sealed class AdminTickOrchestrator(
StorageTickService storageTick,
PipelineProcessingService processing,
PipelineWorkerService worker,
CardsService reminders,
StorageToastPublisher toastPublisher,
SseBroker broker,
PipelinePumpGate pumpGate,
ILogger logger)
{
// SSE-тип события новой карточки (Ruling 5; api.js слушает 'new_card').
private const string NewCardEventType = "new_card";
// SSE-тип события «выстрелившего» напоминания «Отложено» (Ruling 8, api-map §2: {id,title,containerId}).
private const string ReminderDueEventType = "reminder_due";
// Ключи счётчиков pump в pipeline-словаре ответа (1:1 со словарём _pump_unlocked python L921).
private static readonly string[] PipelineCounterKeys =
[
"staged", "rulesStored", "mlStored", "mlDrop", "typeDrop", "aiStored", "aiDrop", "aiFail", "noBudget",
];
///
/// Выполняет один ручной тик тенанта: правила хранения + очистка отсева + напоминания + pump + SSE-публикации.
///
/// Тенант-получатель (сессия запроса; канал SSE-публикаций).
/// Токен отмены запроса.
/// Ответ {storage, reminders, pipeline, queue}; сбой проверки напоминаний/pump не выбрасывается наружу.
public async Task TickAsync(Guid tenantId, CancellationToken ct)
{
// (1) Тик правил хранения канбана (как этап 3; leads.py tick_storage L454–484).
StorageTickStatsDto storage = await storageTick.TickAsync(ct);
// (2) Очистка отсева пайплайна: записи старше 3 суток — безвозвратно (tick_storage L485–486); счётчик
// вливается в storage.purgedRejected (Ruling 9: ответ тика объединяет статистику, L488–493).
int purgedRejected = await processing.PurgeExpiredAsync(ct);
StorageTickStatsDto mergedStorage = storage with { PurgedRejected = purgedRejected };
// (3) Тосты статистики — до pump, как в прототипе (L333): очистка отсева видна, даже если pump упадёт.
toastPublisher.PublishTickToasts(tenantId, mergedStorage);
// (4) Проверка наступивших напоминаний «Отложено» (admin_tick L334 → check_reminders L264–282; план
// Task 11, Ruling 3): CheckDueAsync помечает due-строки fired и возвращает их {id,title,stage}. Сбой
// проверки НЕ роняет тик: лог + reminders ответа пуст (очередь/хранение продолжают работать).
IReadOnlyList dueReminders;
try
{
dueReminders = await reminders.CheckDueRemindersAsync(ct);
}
catch (OperationCanceledException)
{
// Запрос отменён — прерываем тик штатно (не «сбой проверки напоминаний»).
throw;
}
catch (Exception exception)
{
logger.LogWarning(exception, "POST /api/admin/tick: проверка напоминаний не удалась — reminders ответа пуст");
dueReminders = Array.Empty();
}
// SSE reminder_due по каждому «выстрелившему» напоминанию (Ruling 8: событие {id,title,stage}, toast НЕ
// шлём — у фронта модалка ReminderNotice; без подписчиков канала публикация — no-op). После MarkFired
// (внутри CheckDueAsync), как прототип L277–281 — публикуются уже «сработавшие» записи.
foreach (CardReminderDueDto due in dueReminders)
{
broker.Publish(tenantId, ReminderDueEventType, due);
}
// (5) Один проход pump под гейтом тенанта; сбой не роняет тик: pipeline={}, очередь дождётся
// следующего тика/фонового цикла (Ruling 10; Task 11 — гейт общий с фоновым циклом).
PipelinePumpResult? pump = await PumpOnceSafelyAsync(tenantId, ct);
// (6) SSE new_card по карточкам, созданным проходом (Ruling 8/9; без подписчиков — no-op).
if (pump is not null)
{
foreach (CardDto card in pump.CreatedCards)
{
broker.Publish(tenantId, NewCardEventType, card);
}
}
// (7) Строк очереди после pump (queue_len L337: total = new + ai).
QueueCountsDto queueCounts = await processing.QueueCountsAsync(ct);
return new AdminTickResultDto(
mergedStorage,
dueReminders,
pump is null ? new Dictionary() : ToPipelineWireDict(pump),
queueCounts.Total);
}
// Один проход воркера под гейтом тенанта с изоляцией сбоя: исключения pump не роняют тик (Task 10).
// tenantId: Тенант тика (ключ гейта, общего с фоновым циклом Task 11).
// ct: Токен отмены запроса.
// Возвращает: Результат прохода либо null — гейт занят другим воркером/pump упал (pipeline ответа пуст).
private async Task PumpOnceSafelyAsync(Guid tenantId, CancellationToken ct)
{
// Общий воркер-гейт (Ruling 8, аналог asyncio.Lock pipeline.py L40): admin/tick и фоновый цикл не
// разбирают очередь тенанта одновременно. Гейт занят (фоновый цикл уже pump'ит) — проход пропускаем,
// как прототип при занятом локе (L901–902): pipeline={}, очередь дождётся следующего срабатывания.
if (!pumpGate.TryEnter(tenantId))
{
return null;
}
try
{
return await worker.PumpOnceAsync(ct);
}
catch (OperationCanceledException)
{
// Запрос отменён — прерываем тик штатно (не «сбой pump»).
throw;
}
catch (Exception exception)
{
logger.LogWarning(exception, "POST /api/admin/tick: проход pump не удался — тик возвращает storage без pipeline");
return null;
}
finally
{
pumpGate.Exit(tenantId);
}
}
// Счётчики результата pump → pipeline-словарь ответа (9 ключей словаря python L921; карточки в
// wire не выходят — они ушли отдельными SSE new_card).
// pump: Результат успешного прохода pump.
// Возвращает: Словарь счётчиков в wire-порядке прототипа.
private static IReadOnlyDictionary ToPipelineWireDict(PipelinePumpResult pump)
{
int[] counters =
[
pump.Staged, pump.RulesStored, pump.MlStored, pump.MlDrop, pump.TypeDrop,
pump.AiStored, pump.AiDrop, pump.AiFail, pump.NoBudget,
];
var wire = new Dictionary(PipelineCounterKeys.Length, StringComparer.Ordinal);
for (int index = 0; index < PipelineCounterKeys.Length; index++)
{
wire[PipelineCounterKeys[index]] = counters[index];
}
return wire;
}
}