using Deal.Infrastructure.Integrations.Abstractions;
using Deal.Infrastructure.Integrations.Exceptions;
using Deal.Infrastructure.Integrations.Extensions;
using Deal.Infrastructure.Integrations.Models;
using Deal.Infrastructure.Integrations.Options;
using Deal.Infrastructure.Integrations.Services;
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.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;
using Deal.Api.Services;
using Deal.Api.Dtos;
namespace Deal.Api.Hosting;
///
/// Фоновый флашер очереди обучения ML — выгрузка MlOutbox в ml-service батчами (план Task 16, Ruling 6).
///
///
/// Аналог _ml_sync_loop python-прототипа и flush_outbox (ml_client.py L56–82): каждые 10 с
/// обходит ВСЕ тенанты реестра и в собственном scope с ITenantContext.SetTenant (эталон
/// PipelineWorkerScheduler/StorageTickScheduler) отправляет накопленное обучение RPC TrainBatch порциями по
/// 10 строк, ≤100 за цикл. Строки удаляются ТОЛЬКО после успешного батча (python L78–79); при недоступности
/// ml-service порция остаётся и уходит в следующий цикл (ретрай на каждом тике, «строки остаются» — Ruling 6).
///
/// Регистрируется в Deal.Api только при Services:Ml:UseLocal=false (gRPC-режим): Local-режиму
/// ml-service не нужен — очередь копится (этап 3), а при «поднятом сервисе» флашер выгружает её сразу.
/// Первый проход — сразу после старта (как PipelineWorkerScheduler); перекрывающиеся проходы исключены
/// in-flight guard (Interlocked). Ошибки логируются и наружу не выбрасываются (флаш одного тенанта не
/// валит цикл — остальные тенанты обрабатываются); при остановке хоста таймер останавливается и текущий
/// проход отменяется (graceful). Пустая очередь — тихий no-op.
///
///
public sealed class MlOutboxFlushScheduler : IHostedService
{
///
/// Период циклов выгрузки — 10 с (Ruling 6).
///
public const int FlushPeriodSeconds = 10;
///
/// Размер порции за один TrainBatch — 10 строк (ml_client.flush_outbox L63: chunk=10).
///
public const int BatchSize = 10;
///
/// Потолок выгрузки за один цикл тенанта — 100 строк (ml_client.flush_outbox L56: batch=100).
///
public const int MaxPerCycle = 100;
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
// Отмена при остановке хоста: прерывает текущий проход (EF-запросы наблюдают токен).
private readonly CancellationTokenSource _shutdownCts = new();
private Timer? _timer;
private Task? _currentIteration;
private int _iterationInProgress;
///
/// Создаёт планировщик фоновой выгрузки очереди обучения ML.
///
/// Фабрика scope: проход цикла и флаш каждого тенанта — в собственных scope.
/// Логгер сбоев цикла.
public MlOutboxFlushScheduler(IServiceScopeFactory scopeFactory, ILogger logger)
{
ArgumentNullException.ThrowIfNull(scopeFactory);
ArgumentNullException.ThrowIfNull(logger);
_scopeFactory = scopeFactory;
_logger = logger;
}
///
public Task StartAsync(CancellationToken ct)
{
// Первый проход — сразу после старта (как PipelineWorkerScheduler), далее каждые 10 с.
_timer = new Timer(
static state => ((MlOutboxFlushScheduler)state!).RunIteration(),
this,
TimeSpan.Zero,
TimeSpan.FromSeconds(FlushPeriodSeconds));
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)
{
// Лимит остановки истёк — хост продолжает остановку; ошибки прохода проглочены внутри
// FlushCycleCoreAsync, незавершённый проход безопасно завершится на отменённом токене.
}
}
}
///
/// Один проход цикла: список тенантов реестра и выгрузка каждого (no-op, если проход уже идёт).
///
/// Публичен как точка запуска прохода для unit-тестов (тайминги цикла не тестируются);
/// таймер вызывает этот же метод. Ошибки и отмена токена наружу не выбрасываются: сбои логируются
/// (цикл живёт), отмена по токену останова завершает проход штатно.
/// Токен отмены прохода (в проде — токен остановки хоста).
/// Задача прохода (завершается без исключений).
public Task RunCycleAsync(CancellationToken ct)
{
if (Interlocked.CompareExchange(ref _iterationInProgress, 1, 0) != 0)
{
return Task.CompletedTask;
}
Task iteration = FlushCycleCoreAsync(ct);
Volatile.Write(ref _currentIteration, iteration);
return iteration;
}
// Запускает проход из callback таймера (guard — внутри RunCycleAsync).
private void RunIteration()
{
_ = RunCycleAsync(_shutdownCts.Token);
}
// Тело прохода: выгрузка очереди каждого тенанта реестра; guard сбрасывается в finally.
// ct: Токен отмены (остановка хоста).
private async Task FlushCycleCoreAsync(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 FlushTenantAsync(tenant, ct);
}
}
catch (OperationCanceledException)
{
// Остановка хоста: проход прерван по токену — штатный выход, не ошибка.
}
catch (Exception exception)
{
// Сбой всего прохода (реестр недоступен и т.п.): логируем, цикл продолжит со следующего тика.
_logger.LogError(exception, "Флашер ML-outbox: проход по тенантам не удался");
}
finally
{
Interlocked.Exchange(ref _iterationInProgress, 0);
}
}
// Выгрузка очереди одного тенанта в собственном scope: SetTenant → порции по 10 до ≤100/цикл.
// Строки удаляются только после успешного TrainBatch (Ruling 6); сбой батча — порция остаётся,
// цикл тенанта завершается (следующая попытка — следующий тик). Сбой хранилища тенанта не валит проход:
// ошибка логируется, остальные тенанты обрабатываются; отмена (OCE) пробрасывается наверх.
// tenant: Тенант реестра (Id в формате Guid; схема — tenant_<N>).
// ct: Токен отмены прохода.
private async Task FlushTenantAsync(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: IMlLearningStore строится от схемы текущего тенанта (как в цикле pump).
IMlLearningStore learningStore = tenantScope.ServiceProvider.GetRequiredService();
IMlTrainClient trainClient = tenantScope.ServiceProvider.GetRequiredService();
int total = 0;
while (total < MaxPerCycle)
{
IReadOnlyList rows = await learningStore.TakeOutboxBatchAsync(BatchSize, ct);
if (rows.Count == 0)
{
break;
}
try
{
await trainClient.TrainBatchAsync(rows, ct);
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception exception)
{
// Недоступность/сбой ml-service: порция остаётся в очереди (python L75–77), следующая
// попытка — на следующем тике; флашер не роняет проход цикла.
_logger.LogWarning(
exception,
"Флашер ML-outbox: отправка {RowCount} строк тенанта {TenantId} не удалась — строки остались",
rows.Count,
tenant.Id);
break;
}
// Удаление только после успеха (python L78–79): отправленные строки больше не нужны.
await learningStore.DeleteOutboxAsync(rows.Select(row => row.Id).ToList(), ct);
total += rows.Count;
}
if (total > 0)
{
_logger.LogInformation("Флашер ML-outbox: выгружено {RowCount} строк тенанта {TenantId}", total, tenant.Id);
}
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception exception)
{
_logger.LogWarning(exception, "Флашер ML-outbox: выгрузка тенанта {TenantId} не удалась", tenant.Id);
}
finally
{
tenantContext.Reset();
}
}
}