75 lines
3.2 KiB
C#
75 lines
3.2 KiB
C#
using Microsoft.Extensions.Options;
|
|
using PlotLine.Models;
|
|
|
|
namespace PlotLine.Services;
|
|
|
|
public sealed class PersistedStoryIntelligenceWorker(
|
|
IServiceScopeFactory scopeFactory,
|
|
IOptions<StoryIntelligenceOptions> options,
|
|
ILogger<PersistedStoryIntelligenceWorker> logger) : BackgroundService
|
|
{
|
|
private readonly StoryIntelligenceOptions settings = options.Value;
|
|
|
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
|
{
|
|
var workerCount = Math.Clamp(settings.MaxGlobalConcurrentAiRequests, 1, 24);
|
|
logger.LogInformation(
|
|
"Starting {WorkerCount} persisted Story Intelligence worker loop(s). MaxPerBook={MaxPerBook} MaxPerUser={MaxPerUser} PollSeconds={PollSeconds} ClaimLeaseMinutes={ClaimLeaseMinutes}",
|
|
workerCount,
|
|
Math.Max(1, settings.MaxConcurrentAiRequestsPerBook),
|
|
Math.Max(1, settings.MaxConcurrentAiRequestsPerUser),
|
|
Math.Max(1, settings.QueuePollingIntervalSeconds),
|
|
Math.Max(15, settings.ClaimLeaseMinutes));
|
|
|
|
var workers = Enumerable.Range(1, workerCount)
|
|
.Select(workerNumber => RunWorkerLoopAsync(workerNumber, stoppingToken))
|
|
.ToArray();
|
|
await Task.WhenAll(workers);
|
|
}
|
|
|
|
private async Task RunWorkerLoopAsync(int workerNumber, CancellationToken stoppingToken)
|
|
{
|
|
var pollDelay = TimeSpan.FromSeconds(Math.Max(1, settings.QueuePollingIntervalSeconds));
|
|
while (!stoppingToken.IsCancellationRequested)
|
|
{
|
|
var processedRun = false;
|
|
try
|
|
{
|
|
using var scope = scopeFactory.CreateScope();
|
|
var runner = scope.ServiceProvider.GetRequiredService<IPersistedStoryIntelligenceRunner>();
|
|
processedRun = await runner.ProcessNextAsync(stoppingToken);
|
|
}
|
|
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
|
{
|
|
return;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
logger.LogError(ex, "Persisted Story Intelligence worker {WorkerNumber} failed while processing the next queued run.", workerNumber);
|
|
}
|
|
|
|
if (!processedRun)
|
|
{
|
|
await RecoverReadySceneAnalysisAsync(stoppingToken);
|
|
await Task.Delay(pollDelay, stoppingToken);
|
|
}
|
|
}
|
|
}
|
|
|
|
private async Task RecoverReadySceneAnalysisAsync(CancellationToken stoppingToken)
|
|
{
|
|
using var scope = scopeFactory.CreateScope();
|
|
var pipelineState = scope.ServiceProvider.GetRequiredService<IStoryIntelligencePipelineStateService>();
|
|
var notifier = scope.ServiceProvider.GetRequiredService<IStoryIntelligenceProgressNotifier>();
|
|
var completionNotifications = scope.ServiceProvider.GetRequiredService<IStoryIntelligenceCompletionNotificationService>();
|
|
var completion = await pipelineState.RecoverNextReadySceneAnalysisAsync();
|
|
if (completion is not { CompletedNow: true })
|
|
{
|
|
return;
|
|
}
|
|
|
|
await notifier.PublishGlobalAsync(completion.UserID);
|
|
await completionNotifications.NotifyIfImportTerminalAsync(completion.BookID, completion.UserID, stoppingToken);
|
|
}
|
|
}
|