using System.Threading.Channels; using CareFix.Api.Options; using Microsoft.Extensions.Options; namespace CareFix.Api.Ai; public sealed record AiJob(int TicketId, int? UserId, string Text); public sealed class AiQueue { private readonly Channel _channel = Channel.CreateUnbounded(); public ValueTask EnqueueAsync(AiJob job) => _channel.Writer.WriteAsync(job); public ChannelReader Reader => _channel.Reader; } /// Runs AI jobs in the background so the console never waits on a long HTTP request. public sealed class AiWorker(AiQueue queue, AiOrchestrator orchestrator, IOptions options, ILogger log) : BackgroundService { protected override Task ExecuteAsync(CancellationToken stoppingToken) { var workers = Math.Clamp(options.Value.AiWorkers, 1, 16); return Task.WhenAll(Enumerable.Range(0, workers).Select(_ => Task.Run(async () => { await foreach (var job in queue.Reader.ReadAllAsync(stoppingToken)) { try { await orchestrator.RunAsync(job, stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { log.LogError(ex, "AI job failed for ticket {TicketId}", job.TicketId); } } }, stoppingToken))); } }