PlotDirector/PlotLine/Services/PersistedStoryIntelligenceWorker.cs

92 lines
3.9 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 materialisation = scope.ServiceProvider.GetRequiredService<IStoryIntelligenceMaterialisationService>();
var notifier = scope.ServiceProvider.GetRequiredService<IStoryIntelligenceProgressNotifier>();
var completionNotifications = scope.ServiceProvider.GetRequiredService<IStoryIntelligenceCompletionNotificationService>();
var candidate = await pipelineState.FindNextReadySceneAnalysisAsync();
if (candidate is null)
{
return;
}
var materialised = await materialisation.MaterialiseReadyBookAsync(candidate.BookID, candidate.UserID, stoppingToken);
if (!materialised.IsReady)
{
logger.LogInformation(
"Story Intelligence recovery finalisation deferred for BookID={BookID}. Reason={Reason}",
candidate.BookID,
materialised.Reason);
return;
}
var completion = await pipelineState.CompleteSceneAnalysisIfReadyAsync(candidate.BookID, candidate.UserID);
if (completion is not { CompletedNow: true })
{
return;
}
await notifier.PublishGlobalAsync(completion.UserID);
await completionNotifications.NotifyIfImportTerminalAsync(completion.BookID, completion.UserID, stoppingToken);
}
}