using AccountsDAL.CustomCode.Warehouse; using Dapper; using GB5Shared.Connection; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.Telemetry; using Microsoft.AspNetCore.Mvc; using Microsoft.Data.SqlClient; namespace AccountsSL.Subscriptions { // Dapr cron binding handler — triggered by components/warehouse-drain-cron.yaml. // Drains TWAREHOUSECHANGEQUEUE for every POSTINGMODE=1 fact across all active tenants. // // Concurrency: the static SemaphoreSlim(1,1) mirrors Quartz's [DisallowConcurrentExecution] // at the process level. If Dapr fires a second trigger while one run is still in flight, // the handler returns 202 Accepted (2xx so Dapr does NOT retry) and logs a skip notice. // // The tenant-iteration shape is identical to the former WarehouseDeltaDrainJob. It lives // here rather than in a separate service because it's infrastructure plumbing (not business // logic): IServiceScopeFactory + per-tenant IWarehouseChangeQueueBLL scope is the right // tool here, and it doesn't need to be testable independently. [ApiController] public class WarehouseDrainCronHandler : ControllerBase { private static readonly SemaphoreSlim Gate = new(1, 1); private const string TenantQuery = @" SELECT SERVERCONFIG1.CLIENTID AS ClientId, SERVERCONFIG1.DATABASENAME AS DatabaseName, SERVERCONFIG1.DATABASETYPE AS DbType, SERVERCONFIG1.CONNECTIONNAME AS ConnectionName FROM MSERVERCONFIG SERVERCONFIG1 JOIN MSERVER SERVER1 ON SERVERCONFIG1.SERVERID = SERVER1.SERVERID WHERE SERVERCONFIG1.STATUS = 1 AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI'"; private readonly IServiceScopeFactory _scopeFactory; private readonly IApplicationConnection _appConnection; private readonly ILogger _logger; public WarehouseDrainCronHandler( IServiceScopeFactory scopeFactory, IApplicationConnection appConnection, ILogger logger) { _scopeFactory = scopeFactory; _appConnection = appConnection; _logger = logger; } [HttpPost("/Warehouse/DrainCron")] public async Task HandleAsync(CancellationToken ct) { // Non-blocking tryAcquire — 202 tells Dapr "received OK, don't retry" while // making clear this trigger was a no-op because one is already running. if (!await Gate.WaitAsync(0, ct)) { _logger.LogInformation("WarehouseDrainCron: previous run still active — trigger skipped."); return Accepted(); } using var activity = GB5Trace.BeginSection("warehouse-drain-cron"); try { var tenants = await LoadTenantsAsync(ct).ConfigureAwait(false); _logger.LogInformation("WarehouseDrainCron: draining for {Count} tenant(s)", tenants.Count); foreach (var tenant in tenants) { ct.ThrowIfCancellationRequested(); await DrainTenantAsync(tenant, ct).ConfigureAwait(false); } return Ok(); } catch (OperationCanceledException) { _logger.LogWarning("WarehouseDrainCron: cancelled mid-run."); return StatusCode(500); // non-2xx so Dapr retries after next schedule tick } catch (Exception ex) { GB5Trace.MarkFailed("warehouse-drain-cron-failed", ex); _logger.LogError(ex, "WarehouseDrainCron: unhandled error."); return StatusCode(500); } finally { Gate.Release(); } } private async Task DrainTenantAsync(ServerConfigDTO tenant, CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var queueDal = scope.ServiceProvider.GetRequiredService(); var queueBll = scope.ServiceProvider.GetRequiredService(); var login = new LoginDTO { ClientId = tenant.ClientId, UserId = -1, DatabaseName = tenant.DatabaseName, ConnectionDatabaseName = tenant.ConnectionName, IsSchedulerRun = 1 }; List deltaFactIds; try { deltaFactIds = await queueDal.GetDeltaEnabledFactIdsAsync(login, ct).ConfigureAwait(false); } catch (Exception ex) { _logger.LogError(ex, "WarehouseDrainCron: failed to load delta-enabled facts for tenant {ClientId}", tenant.ClientId); return; } if (deltaFactIds.Count == 0) return; foreach (var factId in deltaFactIds) { ct.ThrowIfCancellationRequested(); try { var result = await queueBll.ProcessDeltaForFactAsync(factId, login, ct).ConfigureAwait(false); if (result.SucceededOuIds.Count > 0 || result.FailedOuIds.Count > 0) _logger.LogInformation( "WarehouseDrainCron: tenant {ClientId} factId {FactId} — {Ok} ok, {Fail} failed", tenant.ClientId, factId, result.SucceededOuIds.Count, result.FailedOuIds.Count); } catch (Exception ex) { GB5Trace.MarkFailed("warehouse-drain-cron-fact-failed", ex); _logger.LogError(ex, "WarehouseDrainCron: tenant {ClientId} factId {FactId} failed", tenant.ClientId, factId); } } } private async Task> LoadTenantsAsync(CancellationToken ct) { var systemConn = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false); using var conn = new SqlConnection(systemConn); var rows = await conn.QueryAsync(TenantQuery).ConfigureAwait(false); return rows.ToList(); } } }