using System.Threading.Channels; using Microsoft.Extensions.Options; using UnrealDemoScanner.Web.Infrastructure; using UnrealDemoScanner.Web.Models; namespace UnrealDemoScanner.Web.Services; public sealed class ScanQueueService : BackgroundService, IScanQueueService { private readonly Channel _channel; private readonly IServiceProvider _serviceProvider; private readonly AppPaths _paths; private readonly ILogger _logger; private readonly ScannerOptions _scannerOptions; public ScanQueueService( IServiceProvider serviceProvider, IOptions scannerOptions, AppPaths paths, ILogger logger) { _serviceProvider = serviceProvider; _paths = paths; _logger = logger; _scannerOptions = scannerOptions.Value; _channel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = false, SingleWriter = false, AllowSynchronousContinuations = false }); } public ValueTask EnqueueAsync(string jobId, CancellationToken cancellationToken = default) { return _channel.Writer.WriteAsync(jobId, cancellationToken); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { var workers = Math.Max(1, _scannerOptions.MaxConcurrentJobs); _logger.LogInformation("Starting scan queue with {WorkerCount} worker(s).", workers); var tasks = Enumerable.Range(0, workers) .Select(worker => WorkerLoop(worker, stoppingToken)) .ToArray(); await Task.WhenAll(tasks); } private async Task WorkerLoop(int workerId, CancellationToken stoppingToken) { while (await _channel.Reader.WaitToReadAsync(stoppingToken)) { while (_channel.Reader.TryRead(out var jobId)) { try { await ProcessJobAsync(workerId, jobId, stoppingToken); } catch (Exception ex) { _logger.LogError(ex, "Worker {WorkerId} failed processing job {JobId}", workerId, jobId); } } } } private async Task ProcessJobAsync(int workerId, string jobId, CancellationToken cancellationToken) { using var scope = _serviceProvider.CreateScope(); var repo = scope.ServiceProvider.GetRequiredService(); var runner = scope.ServiceProvider.GetRequiredService(); var parser = scope.ServiceProvider.GetRequiredService(); var job = await repo.GetJobAsync(jobId, cancellationToken); if (job is null) { _logger.LogWarning("Worker {WorkerId} skipped missing job {JobId}", workerId, jobId); return; } _logger.LogInformation("Worker {WorkerId} starting job {JobId}", workerId, jobId); await repo.MarkRunningAsync(jobId, DateTime.UtcNow, cancellationToken); var resultFolder = Path.Combine(_paths.ResultsRoot, jobId); Directory.CreateDirectory(resultFolder); var outputLogPath = Path.Combine(resultFolder, "scanner-output.log"); try { var run = await runner.RunAsync(job.StoredFilePath, cancellationToken); await File.WriteAllTextAsync(outputLogPath, run.CombinedOutput, cancellationToken); var parsed = parser.Parse(run.CombinedOutput, jobId); var status = run.ExitCode == 0 ? ScanJobStatus.Succeeded : ScanJobStatus.Failed; var error = run.TimedOut ? "Scanner timed out." : (run.ExitCode == 0 ? null : $"Scanner exit code {run.ExitCode}"); await repo.CompleteAsync( id: jobId, status: status, exitCode: run.ExitCode, errorMessage: error, outputLogPath: outputLogPath, durationMs: run.DurationMs, detectedCount: parsed.DetectedCount, warningCount: parsed.WarningCount, infoCount: parsed.InfoCount, findings: parsed.Findings, cancellationToken: cancellationToken); } catch (Exception ex) { _logger.LogError(ex, "Worker {WorkerId} failed running scanner for job {JobId}", workerId, jobId); var errorMessage = ex.Message; var outputText = ex.ToString(); await File.WriteAllTextAsync(outputLogPath, outputText, cancellationToken); await repo.CompleteAsync( id: jobId, status: ScanJobStatus.Failed, exitCode: -1, errorMessage: errorMessage, outputLogPath: outputLogPath, durationMs: 0, detectedCount: 0, warningCount: 0, infoCount: 0, findings: Array.Empty(), cancellationToken: cancellationToken); } } }