using Deal.Api.Events;
using Deal.Modules.Kanban.Application.Abstractions;
using Deal.Modules.Kanban.Application.Extensions;
using Deal.Modules.Kanban.Application.Models;
using Deal.Modules.Kanban.Application.Registrars;
using Deal.Modules.Kanban.Application.Services;
using Deal.Modules.Pipeline.Application.Abstractions;
using Deal.Modules.Pipeline.Application.Models;
using Deal.Modules.Pipeline.Application.Registrars;
using Deal.Modules.Pipeline.Application.Services;
using Deal.Modules.Tenants.Application.Abstractions;
using Deal.Modules.Tenants.Application.Extensions;
using Deal.Modules.Tenants.Application.Models;
using Deal.Modules.Tenants.Application.Registrars;
using Deal.Modules.Tenants.Application.Services;
using Deal.SharedKernel.Tenants;
namespace Deal.Api.Hosting;
///
/// Фоновый цикл правил хранения по всем тенантам (план Tasks 11–12, Rulings 3/8; аналог _storage_loop main.py L43–53).
///
///
/// Каждые 30 с (в прототипе — asyncio.sleep(30), main.py L53) обходит ВСЕ тенанты системного
/// реестра. Проход открывает собственный scope (реестр живёт в публичной схеме — ITenantRepository вне
/// tenant-контекста, паттерн TenantBootstrapService), на каждый тенант — вложенный scope с
/// ITenantContext.SetTenant (эталон SessionMiddleware/TenantBootstrapService) и StorageTickService.TickAsync,
/// после чего — автоочистка отсева пайплайна (, 3 суток;
/// Task 11, Ruling 8/9: фоновый аналог тика AdminTickOrchestrator, tick_storage L485–493), SSE-тосты
/// статистики в канал тенанта (StorageToastPublisher; без подписчиков — no-op, Ruling 5) и проверка
/// наступивших напоминаний «Отложено» (CardsService.CheckDueRemindersAsync + SSE reminder_due, Task 12,
/// Ruling 3/8 — фоновый аналог ветки AdminTickOrchestrator, check_reminders из _storage_loop main.py L49).
/// Порядок тика тенанта 1:1 с _storage_loop main.py L47–53: тик → тосты → напоминания. Как в прототипе,
/// первый проход выполняется сразу после старта (тик до первого sleep), далее — по таймеру. Параллельные
/// проходы исключены in-flight guard (Interlocked, как RatesRefreshScheduler): если проход длится дольше
/// периода, следующее срабатывание таймера пропускается. Ошибки логируются и наружу не выбрасываются
/// (тик одного тенанта не валит проход); при остановке хоста таймер останавливается и текущий проход
/// отменяется (graceful).
///
public sealed class StorageTickScheduler : IHostedService
{
// Период проходов цикла — 30 с, 1:1 с _storage_loop main.py L53 (asyncio.sleep(30)).
private const int TickPeriodSeconds = 30;
// SSE-тип события «выстрелившего» напоминания «Отложено» (Ruling 8, api-map §2: {id,title,containerId}).
private const string ReminderDueEventType = "reminder_due";
private static readonly TimeSpan TickPeriod = TimeSpan.FromSeconds(TickPeriodSeconds);
private readonly IServiceScopeFactory _scopeFactory;
private readonly StorageToastPublisher _toastPublisher;
private readonly SseBroker _broker;
private readonly ILogger _logger;
// Отмена при остановке хоста: прерывает текущий проход (EF-запросы тика наблюдают токен).
private readonly CancellationTokenSource _shutdownCts = new();
private Timer? _timer;
private Task? _currentIteration;
private int _iterationInProgress;
///
/// Создаёт планировщик фонового цикла правил хранения.
///
/// Фабрика scope: проход цикла и тик каждого тенанта — в собственных scope.
/// Публикатор SSE-тостов статистики тика (общий с POST /api/admin/tick).
/// SSE-брокер каналов тенантов: reminder_due «выстреливших» напоминаний в канал тенанта
/// (как StorageToastPublisher; без подписчиков публикация — no-op, Ruling 5).
/// Логгер ошибок цикла.
public StorageTickScheduler(
IServiceScopeFactory scopeFactory,
StorageToastPublisher toastPublisher,
SseBroker broker,
ILogger logger)
{
ArgumentNullException.ThrowIfNull(scopeFactory);
ArgumentNullException.ThrowIfNull(toastPublisher);
ArgumentNullException.ThrowIfNull(broker);
ArgumentNullException.ThrowIfNull(logger);
_scopeFactory = scopeFactory;
_toastPublisher = toastPublisher;
_broker = broker;
_logger = logger;
}
///
public Task StartAsync(CancellationToken ct)
{
// Первый проход — сразу после старта (в прототипе тик выполняется до первого sleep); далее каждые 30 с.
_timer = new Timer(
static state => ((StorageTickScheduler)state!).RunIteration(),
this,
TimeSpan.Zero,
TickPeriod);
return Task.CompletedTask;
}
///
public async Task StopAsync(CancellationToken ct)
{
// Новые проходы не запускаем; текущий отменяем и ждём его завершения — не дольше лимита
// остановки хоста (HostOptions.ShutdownTimeout).
_timer?.Change(Timeout.InfiniteTimeSpan, Timeout.InfiniteTimeSpan);
_timer?.Dispose();
_timer = null;
_shutdownCts.Cancel();
Task? iteration = Volatile.Read(ref _currentIteration);
if (iteration is not null)
{
try
{
await iteration.WaitAsync(ct);
}
catch (OperationCanceledException)
{
// Лимит остановки истёк — хост продолжает остановку; ошибки прохода проглочены внутри
// RunCycleAsync, незавершённый проход безопасно завершится на отменённом токене.
}
}
}
///
/// Один проход цикла: список тенантов реестра и тик каждого (no-op, если проход уже идёт).
///
/// Публичен как точка запуска прохода для unit-тестов (тайминги цикла не тестируются) и
/// ручного вызова при отладке; таймер вызывает этот же метод. Ошибки и отмена токена наружу не
/// выбрасываются: сбои логируются (цикл живёт), отмена по токену останова завершает проход штатно.
/// Токен отмены прохода (в проде — токен остановки хоста).
/// Задача прохода (завершается без исключений).
public Task RunCycleAsync(CancellationToken ct)
{
if (Interlocked.CompareExchange(ref _iterationInProgress, 1, 0) != 0)
{
return Task.CompletedTask;
}
Task iteration = RunCycleCoreAsync(ct);
Volatile.Write(ref _currentIteration, iteration);
return iteration;
}
// Запускает проход из callback таймера (guard — внутри RunCycleAsync).
private void RunIteration()
{
_ = RunCycleAsync(_shutdownCts.Token);
}
// Тело прохода: тик каждого тенанта реестра; guard сбрасывается в finally.
// ct: Токен отмены (остановка хоста).
private async Task RunCycleCoreAsync(CancellationToken ct)
{
try
{
await using AsyncServiceScope cycleScope = _scopeFactory.CreateAsyncScope();
ITenantRepository tenantRepository = cycleScope.ServiceProvider.GetRequiredService();
IReadOnlyList tenants = await tenantRepository.ListAsync(ct);
foreach (TenantRecordDto tenant in tenants)
{
await TickTenantAsync(tenant, ct);
}
}
catch (OperationCanceledException)
{
// Остановка хоста: проход прерван по токену — штатный выход, не ошибка.
}
catch (Exception exception)
{
// Сбой всего прохода (реестр недоступен и т.п.): логируем, цикл продолжит со следующего тика.
_logger.LogError(exception, "Цикл правил хранения: проход по тенантам не удался");
}
finally
{
Interlocked.Exchange(ref _iterationInProgress, 0);
}
}
// Тик одного тенанта в собственном scope: SetTenant → Kanban-тик → purge отсева → SSE-тосты →
// проверка напоминаний + SSE reminder_due; Reset в finally (1:1 с _storage_loop main.py L47–53: тик →
// тосты → напоминания).
// Контекст AsyncLocal сбрасывается в finally, чтобы не переживать scope тенанта (как
// SessionMiddleware). Тик одного тенанта не валит проход: ошибка ветки/тенанта логируется, остальные
// тенанты обрабатываются; отмена (OCE) пробрасывается наверх — проход завершается.
// tenant: Тенант реестра (Id в формате Guid; схема — tenant_<N>).
// ct: Токен отмены прохода.
private async Task TickTenantAsync(TenantRecordDto tenant, CancellationToken ct)
{
await using AsyncServiceScope tenantScope = _scopeFactory.CreateAsyncScope();
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService();
try
{
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
// Resolve ПОСЛЕ SetTenant: TenantDbContext (и его адаптеры) строятся от схемы текущего тенанта.
StorageTickService tickService = tenantScope.ServiceProvider.GetRequiredService();
StorageTickStatsDto stats = await tickService.TickAsync(ct);
// Автоочистка отсева пайплайна: записи старше 3 суток — безвозвратно (Task 11, Ruling 8/9;
// tick_storage L485–493). Счётчик вливается в storage.purgedRejected — как ручной тик
// (AdminTickOrchestrator), тост «Отсев очищен: N записей (3 дн.)» публикуется этой же веткой.
PipelineProcessingService processing = tenantScope.ServiceProvider.GetRequiredService();
int purgedRejected = await processing.PurgeExpiredAsync(ct);
StorageTickStatsDto mergedStats = stats with { PurgedRejected = purgedRejected };
// Тосты — в канал тенанта (публикация из Api-слоя, Ruling 5; без подписчиков — no-op).
_toastPublisher.PublishTickToasts(tenant.Id, mergedStats);
// Проверка наступивших напоминаний «Отложено» (Task 12, Ruling 3/8; proj_svc.check_reminders из
// _storage_loop main.py L49): CheckDueRemindersAsync помечает due-строки hold-карточек fired и
// возвращает их {id,title,containerId}. Сбой проверки НЕ роняет тик тенанта/проход: лог-предупреждение,
// остальные тенанты обрабатываются (паттерн ветки AdminTickOrchestrator). SSE reminder_due по каждой
// записи — в канал тенанта (Ruling 8: публикации только из Api; toast НЕ шлём, без подписчиков — no-op).
IReadOnlyList dueReminders;
try
{
CardsService reminders = tenantScope.ServiceProvider.GetRequiredService();
dueReminders = await reminders.CheckDueRemindersAsync(ct);
}
catch (OperationCanceledException)
{
// Остановка хоста — прерываем проход штатно (не «сбой проверки напоминаний»).
throw;
}
catch (Exception exception)
{
_logger.LogWarning(exception, "Цикл правил хранения: проверка напоминаний тенанта {TenantId} не удалась", tenant.Id);
dueReminders = Array.Empty();
}
foreach (CardReminderDueDto due in dueReminders)
{
_broker.Publish(tenant.Id, ReminderDueEventType, due);
}
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception exception)
{
_logger.LogWarning(exception, "Цикл правил хранения: тик тенанта {TenantId} не удался", tenant.Id);
}
finally
{
tenantContext.Reset();
}
}
}