< Summary

Line coverage
94%
Covered lines: 72
Uncovered lines: 4
Coverable lines: 76
Total lines: 454
Line coverage: 94.7%
Branch coverage
100%
Covered branches: 12
Total branches: 12
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
File 1: .cctor()100%210%
File 2: .ctor(...)100%11100%
File 2: RunAsync()100%11100%
File 2: ProduceAsync()100%22100%
File 2: EnqueueDirectoryItemsAsync()100%88100%
File 2: ConsumeAsync()100%22100%

File(s)

/home/runner/work/anichron/anichron/src/.artifacts/obj/Anichron.Worker/debug/Microsoft.Extensions.Logging.Generators/Microsoft.Extensions.Logging.Generators.LoggerMessageGenerator/LoggerMessage.g.cs

File '/home/runner/work/anichron/anichron/src/.artifacts/obj/Anichron.Worker/debug/Microsoft.Extensions.Logging.Generators/Microsoft.Extensions.Logging.Generators.LoggerMessageGenerator/LoggerMessage.g.cs' does not exist (any more).

/home/runner/work/anichron/anichron/src/Anichron.Worker/Crawling/FileIngestionPipeline.cs

#LineLine coverage
 1using Anichron.Core;
 2using Anichron.Core.Domain;
 3using Anichron.Worker.Ingestion;
 4using Anichron.Worker.Ingestion.Pipeline;
 5using Anichron.Worker.Settings;
 6using Microsoft.Extensions.Options;
 7using System.IO.Abstractions;
 8using System.Threading.Channels;
 9
 10namespace Anichron.Worker.Crawling;
 11
 12internal interface IFileIngestionPipeline
 13{
 14    Task RunAsync(UserStorageConfig config, CancellationToken ct);
 15}
 16
 1317internal sealed partial class FileIngestionPipeline(
 1318    IServiceScopeFactory scopeFactory,
 1319    IFileSystem fileSystem,
 1320    ILivePhotoLinker livePhotoLinker,
 1321    IOptions<WorkerSettings> workerOptions,
 1322    IGuidFactory guidFactory,
 1323    ILogger<FileIngestionPipeline> logger) : IFileIngestionPipeline
 24{
 1325    private readonly WorkerSettings settings = workerOptions.Value;
 26
 27    public async Task RunAsync(UserStorageConfig config, CancellationToken ct)
 1328    {
 1329        var channel = Channel.CreateBounded<IngestionItem>(
 1330            new BoundedChannelOptions(settings.MaxConcurrentFiles * 2)
 1331            {
 1332                SingleWriter = true,
 1333                FullMode = BoundedChannelFullMode.Wait,
 1334            });
 35
 1336        var producer = ProduceAsync(config, channel.Writer, ct);
 1337        var consumers = Enumerable
 1338            .Range(0, settings.MaxConcurrentFiles)
 1339            .Select(workerIndex => ConsumeAsync(config, channel.Reader, workerIndex, ct))
 1340            .ToArray();
 41
 1342        await Task.WhenAll([producer, .. consumers]);
 1143    }
 44
 45    private async Task ProduceAsync(
 46        UserStorageConfig config,
 47        ChannelWriter<IngestionItem> writer,
 48        CancellationToken ct)
 1349    {
 50        try
 1351        {
 1352            var filesByDirectory = fileSystem.Directory
 1353                .EnumerateFiles(config.RootPath, "*", SearchOption.AllDirectories)
 1354                .GroupBy(path => fileSystem.Path.GetDirectoryName(path) ?? string.Empty);
 55
 6056            foreach (var directoryGroup in filesByDirectory)
 1257                await EnqueueDirectoryItemsAsync([.. directoryGroup], config.RootPath, writer, ct);
 1258        }
 59        finally
 1360        {
 1361            writer.TryComplete();
 1362        }
 1263    }
 64
 65    private async Task EnqueueDirectoryItemsAsync(
 66        IList<string> filesInDirectory,
 67        string rootPath,
 68        ChannelWriter<IngestionItem> writer,
 69        CancellationToken ct)
 1270    {
 1271        var linkResult = livePhotoLinker.Link(filesInDirectory, rootPath);
 72
 4673        foreach (var item in linkResult.Items)
 574            await writer.WriteAsync(item, ct);
 75
 7476        foreach (var filePath in filesInDirectory)
 1977        {
 1978            if (linkResult.ClaimedPaths.Contains(filePath))
 879                continue;
 80
 1181            var mediaType = MediaTypeDetector.Detect(filePath);
 1182            if (mediaType is null)
 283            {
 284                Log.UnsupportedFile(logger, filePath);
 285                continue;
 86            }
 87
 988            var relativePath = fileSystem.Path.GetRelativePath(rootPath, filePath);
 989            await writer.WriteAsync(new SingleFileItem(filePath, relativePath, mediaType.Value), ct);
 990        }
 1291    }
 92
 93    private async Task ConsumeAsync(
 94        UserStorageConfig config,
 95        ChannelReader<IngestionItem> reader,
 96        int workerIndex,
 97        CancellationToken ct)
 2498    {
 9599        await foreach (var item in reader.ReadAllAsync(ct))
 12100        {
 12101            var filename = fileSystem.Path.GetFileName(item.AbsolutePath);
 12102            using var logScope = logger.BeginScope(
 12103                "W{WorkerIndex} {IngestionFile}", workerIndex, filename);
 104            // Scoped services (e.g. DbContext) must not be shared across concurrent files.
 12105            using var diScope = scopeFactory.CreateScope();
 12106            var runner = diScope.ServiceProvider.GetRequiredService<IIngestionPipelineRunner>();
 107            try
 12108            {
 12109                var context = new IngestionContext { Item = item, Config = config, AssetId = guidFactory.NewGuid() };
 12110                await runner.RunAsync(context, ct);
 10111            }
 2112            catch (Exception ex) when (IngestionShutdown.IsInProgress(ex, ct))
 1113            {
 114                // A shutdown is not an item failure: stop promptly instead of draining the queue.
 1115                throw;
 116            }
 1117            catch (Exception ex)
 1118            {
 1119                Log.ItemFailed(logger, item.AbsolutePath, ex);
 1120            }
 11121        }
 23122    }
 123
 124    private static partial class Log
 125    {
 126        [LoggerMessage(Level = LogLevel.Warning, Message = "Skipping unsupported file type: {Path}.")]
 127        public static partial void UnsupportedFile(ILogger logger, string path);
 128
 129        [LoggerMessage(Level = LogLevel.Error, Message = "Failed to ingest {Path}.")]
 130        public static partial void ItemFailed(ILogger logger, string path, Exception ex);
 131    }
 132}