| | | 1 | | using Anichron.Worker.Ingestion.Proxy; |
| | | 2 | | using Anichron.Worker.Settings; |
| | | 3 | | using Microsoft.Extensions.Options; |
| | | 4 | | using System.IO.Abstractions; |
| | | 5 | | |
| | | 6 | | namespace Anichron.Worker.Ingestion.Pipeline; |
| | | 7 | | |
| | | 8 | | internal interface IIngestionPipelineRunner |
| | | 9 | | { |
| | | 10 | | Task RunAsync(IngestionContext context, CancellationToken ct); |
| | | 11 | | } |
| | | 12 | | |
| | 27 | 13 | | internal sealed partial class IngestionPipelineRunner( |
| | 27 | 14 | | IEnumerable<IIngestionMiddleware> middlewares, |
| | 27 | 15 | | IProxyDirectoryStrategy proxyDirectoryStrategy, |
| | 27 | 16 | | IFileSystem fileSystem, |
| | 27 | 17 | | IOptions<WorkerSettings> settings, |
| | 27 | 18 | | ILogger<IngestionPipelineRunner> logger) : IIngestionPipelineRunner |
| | | 19 | | { |
| | 28 | 20 | | private readonly IngestionDelegate pipeline = BuildValidated([.. middlewares], logger); |
| | | 21 | | |
| | | 22 | | public async Task RunAsync(IngestionContext context, CancellationToken ct) |
| | 25 | 23 | | { |
| | | 24 | | try |
| | 25 | 25 | | { |
| | 25 | 26 | | await pipeline(context, ct); |
| | 12 | 27 | | } |
| | 13 | 28 | | catch (Exception ex) when (!IngestionShutdown.IsInProgress(ex, ct)) |
| | 10 | 29 | | { |
| | | 30 | | // The pipeline generates proxies before it writes the asset row, so a failure part-way |
| | | 31 | | // through leaves files on disk that no row will ever reference. This is the one place |
| | | 32 | | // that sees every file written across every step, so cleanup belongs here. |
| | 10 | 33 | | Compensate(context); |
| | 10 | 34 | | throw; |
| | | 35 | | } |
| | 12 | 36 | | } |
| | | 37 | | |
| | | 38 | | private void Compensate(IngestionContext context) |
| | 10 | 39 | | { |
| | 10 | 40 | | var proxyRoot = settings.Value.ProxyPath; |
| | 10 | 41 | | var proxies = context.ProxyFiles; |
| | | 42 | | |
| | 50 | 43 | | foreach (var proxy in proxies) |
| | 10 | 44 | | TryDelete(fileSystem.Path.Combine(proxyRoot, proxy.ProxyPath)); |
| | | 45 | | |
| | | 46 | | // Asking the strategy rather than deriving the directory from the proxy paths keeps the |
| | | 47 | | // path rule in one place, and covers the attempt that failed in its first generator: it |
| | | 48 | | // registered no proxy at all, but the step had already created the directory. |
| | 10 | 49 | | if (context.ContentHash is { } contentHash) |
| | 9 | 50 | | { |
| | 9 | 51 | | var directory = fileSystem.Path.Combine( |
| | 9 | 52 | | proxyRoot, proxyDirectoryStrategy.GetDirectory(context.Config.Id, contentHash)); |
| | 9 | 53 | | DiscardTemporaryFiles(directory); |
| | 9 | 54 | | TryRemoveIfEmpty(directory); |
| | 9 | 55 | | } |
| | | 56 | | |
| | 10 | 57 | | Log.Compensated(logger, proxies.Count, context.Item.RelativePath); |
| | 10 | 58 | | } |
| | | 59 | | |
| | | 60 | | // Deletion runs while the original exception is in flight, so a file that cannot be deleted is |
| | | 61 | | // reported and skipped: one locked file must not strand the others or replace the real error. |
| | | 62 | | private void TryDelete(string absolutePath) |
| | 12 | 63 | | { |
| | | 64 | | try |
| | 12 | 65 | | { |
| | 12 | 66 | | if (fileSystem.File.Exists(absolutePath)) |
| | 12 | 67 | | fileSystem.File.Delete(absolutePath); |
| | 10 | 68 | | } |
| | 2 | 69 | | catch (Exception ex) when (ex is IOException or UnauthorizedAccessException) |
| | 2 | 70 | | { |
| | 2 | 71 | | Log.DeleteFailed(logger, absolutePath, ex); |
| | 2 | 72 | | } |
| | 12 | 73 | | } |
| | | 74 | | |
| | | 75 | | private void DiscardTemporaryFiles(string directory) |
| | 9 | 76 | | { |
| | | 77 | | try |
| | 9 | 78 | | { |
| | 9 | 79 | | var temporaryFiles = fileSystem.Directory.EnumerateFiles( |
| | 9 | 80 | | directory, ProxyStagingWriter.TemporaryFileSearchPattern); |
| | 31 | 81 | | foreach (var temporaryPath in temporaryFiles) |
| | 2 | 82 | | TryDelete(temporaryPath); |
| | 9 | 83 | | } |
| | 0 | 84 | | catch (Exception ex) when (ex is IOException or UnauthorizedAccessException) |
| | 0 | 85 | | { |
| | 0 | 86 | | Log.DeleteFailed(logger, directory, ex); |
| | 0 | 87 | | } |
| | 9 | 88 | | } |
| | | 89 | | |
| | | 90 | | // Only when empty: the directory is keyed by storage config and content hash, so a duplicate |
| | | 91 | | // asset can share it until the unique index in #159 lands, and a recursive delete would |
| | | 92 | | // destroy that sibling's proxies. |
| | | 93 | | private void TryRemoveIfEmpty(string directory) |
| | 9 | 94 | | { |
| | | 95 | | try |
| | 9 | 96 | | { |
| | 9 | 97 | | if (fileSystem.Directory.Exists(directory) |
| | 9 | 98 | | && !fileSystem.Directory.EnumerateFileSystemEntries(directory).Any()) |
| | 6 | 99 | | { |
| | 6 | 100 | | fileSystem.Directory.Delete(directory); |
| | 6 | 101 | | } |
| | 9 | 102 | | } |
| | 0 | 103 | | catch (Exception ex) when (ex is IOException or UnauthorizedAccessException) |
| | 0 | 104 | | { |
| | 0 | 105 | | Log.DeleteFailed(logger, directory, ex); |
| | 0 | 106 | | } |
| | 9 | 107 | | } |
| | | 108 | | |
| | | 109 | | private static IngestionDelegate BuildValidated( |
| | | 110 | | IReadOnlyList<IIngestionMiddleware> middlewares, ILogger logger) |
| | 28 | 111 | | { |
| | 28 | 112 | | Validate(middlewares); |
| | 27 | 113 | | var ordered = middlewares.OrderBy(m => m.Order).ToArray(); |
| | 27 | 114 | | return IngestionPipelineBuilder.Build(ordered, logger); |
| | 27 | 115 | | } |
| | | 116 | | |
| | | 117 | | private static void Validate(IReadOnlyList<IIngestionMiddleware> middlewares) |
| | 28 | 118 | | { |
| | 28 | 119 | | var duplicates = middlewares |
| | 28 | 120 | | .GroupBy(m => m.Order) |
| | 28 | 121 | | .Where(g => g.Count() > 1) |
| | 28 | 122 | | .Select(g => g.Key) |
| | 28 | 123 | | .ToList(); |
| | 28 | 124 | | if (duplicates.Count > 0) |
| | 1 | 125 | | { |
| | 1 | 126 | | throw new InvalidOperationException( |
| | 1 | 127 | | $"Duplicate middleware orders: {string.Join(", ", duplicates)}"); |
| | | 128 | | } |
| | 27 | 129 | | } |
| | | 130 | | |
| | | 131 | | private static partial class Log |
| | | 132 | | { |
| | | 133 | | [LoggerMessage(Level = LogLevel.Information, |
| | | 134 | | Message = "Cleaned up {Count} proxy file(s) after a failed ingestion of {RelativePath}.")] |
| | | 135 | | public static partial void Compensated(ILogger logger, int count, string relativePath); |
| | | 136 | | |
| | | 137 | | [LoggerMessage(Level = LogLevel.Warning, Message = "Could not delete {Path} during cleanup.")] |
| | | 138 | | public static partial void DeleteFailed(ILogger logger, string path, Exception ex); |
| | | 139 | | } |
| | | 140 | | } |