104 lines
4.4 KiB
C#
104 lines
4.4 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);
|
|
processedRun = await ProcessWholeBookPlotSynthesisAsync(stoppingToken);
|
|
}
|
|
|
|
if (!processedRun)
|
|
{
|
|
await Task.Delay(pollDelay, stoppingToken);
|
|
}
|
|
}
|
|
}
|
|
|
|
private async Task<bool> ProcessWholeBookPlotSynthesisAsync(CancellationToken stoppingToken)
|
|
{
|
|
using var scope = scopeFactory.CreateScope();
|
|
var synthesis = scope.ServiceProvider.GetRequiredService<IWholeBookPlotIntelligenceService>();
|
|
return await synthesis.ProcessNextAsync(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);
|
|
}
|
|
}
|