| | | 1 | | using Anichron.Core; |
| | | 2 | | using Anichron.Core.Domain; |
| | | 3 | | using Anichron.Worker.Ingestion.Pipeline; |
| | | 4 | | using Anichron.Worker.Settings; |
| | | 5 | | using Microsoft.Extensions.Options; |
| | | 6 | | using NodaTime; |
| | | 7 | | using System.IO.Abstractions; |
| | | 8 | | |
| | | 9 | | namespace Anichron.Worker.Ingestion.Proxy; |
| | | 10 | | |
| | 47 | 11 | | internal sealed partial class ProxyStagingWriter( |
| | 47 | 12 | | IFileSystem fileSystem, |
| | 47 | 13 | | IOptions<WorkerSettings> settings, |
| | 47 | 14 | | IClock clock, |
| | 47 | 15 | | IGuidFactory guidFactory, |
| | 47 | 16 | | ILogger<ProxyStagingWriter> logger) |
| | | 17 | | { |
| | | 18 | | // A proxy is written under this suffix and renamed into place once complete, so a crash or a |
| | | 19 | | // killed transcode can never leave a truncated file wearing a valid proxy's name. Leftovers |
| | | 20 | | // are recognisable by this suffix, both to an operator and to the orphan sweeper (#174). |
| | | 21 | | internal const string TemporaryFileSuffix = ".tmp"; |
| | | 22 | | |
| | | 23 | | // Naming a leftover and finding one are the same rule, so they live together: a consumer that |
| | | 24 | | // hunts for staging artifacts — compensation here, the orphan sweeper in #174 — must never |
| | | 25 | | // re-express how this writer names them. |
| | | 26 | | internal const string TemporaryFileSearchPattern = "*" + TemporaryFileSuffix; |
| | | 27 | | |
| | | 28 | | internal static string TemporaryPathFor(string proxyAbsolutePath) |
| | 79 | 29 | | => proxyAbsolutePath + TemporaryFileSuffix; |
| | | 30 | | |
| | | 31 | | // The one place a proxy reaches its final name: produceAsync writes the bytes to the temporary |
| | | 32 | | // path it is handed, and only a complete file is renamed into place and registered on the |
| | | 33 | | // context. The rename stays within one directory, so it is atomic. |
| | | 34 | | internal async Task WriteAsync( |
| | | 35 | | IngestionContext context, |
| | | 36 | | string relativePath, |
| | | 37 | | ProxyType proxyType, |
| | | 38 | | Func<string, CancellationToken, Task> produceAsync, |
| | | 39 | | CancellationToken ct) |
| | 78 | 40 | | { |
| | 78 | 41 | | var proxyAbsolutePath = fileSystem.Path.Combine(settings.Value.ProxyPath, relativePath); |
| | 78 | 42 | | var temporaryAbsolutePath = TemporaryPathFor(proxyAbsolutePath); |
| | | 43 | | try |
| | 78 | 44 | | { |
| | 78 | 45 | | if (fileSystem.Path.GetDirectoryName(proxyAbsolutePath) is { } directory) |
| | 78 | 46 | | fileSystem.Directory.CreateDirectory(directory); |
| | | 47 | | |
| | 78 | 48 | | await produceAsync(temporaryAbsolutePath, ct); |
| | 73 | 49 | | var sizeBytes = fileSystem.FileInfo.New(temporaryAbsolutePath).Length; |
| | 73 | 50 | | fileSystem.File.Move(temporaryAbsolutePath, proxyAbsolutePath, overwrite: true); |
| | | 51 | | |
| | | 52 | | // Registered at the moment it lands, so the context always reflects what is on disk |
| | | 53 | | // and compensation can delete precisely the files this attempt wrote. |
| | 73 | 54 | | context.AddProxyFile( |
| | 73 | 55 | | ProxyFileBuilder.Build(context, relativePath, proxyType, sizeBytes, guidFactory, clock)); |
| | 73 | 56 | | Log.ProxyWritten(logger, proxyType, sizeBytes, context.Item.RelativePath); |
| | 73 | 57 | | } |
| | | 58 | | finally |
| | 78 | 59 | | { |
| | 78 | 60 | | DiscardTemporaryFile(temporaryAbsolutePath); |
| | 78 | 61 | | } |
| | 73 | 62 | | } |
| | | 63 | | |
| | | 64 | | // Runs while an exception may already be in flight — including a shutdown cancellation, where |
| | | 65 | | // the half-written file is dropped but completed proxies are kept. A failure to delete must |
| | | 66 | | // never replace the error that caused the cleanup. |
| | | 67 | | private void DiscardTemporaryFile(string temporaryPath) |
| | 78 | 68 | | { |
| | | 69 | | try |
| | 78 | 70 | | { |
| | 78 | 71 | | if (fileSystem.File.Exists(temporaryPath)) |
| | 4 | 72 | | fileSystem.File.Delete(temporaryPath); |
| | 77 | 73 | | } |
| | 1 | 74 | | catch (Exception ex) when (ex is IOException or UnauthorizedAccessException) |
| | 1 | 75 | | { |
| | 1 | 76 | | Log.DiscardFailed(logger, temporaryPath, ex); |
| | 1 | 77 | | } |
| | 78 | 78 | | } |
| | | 79 | | |
| | | 80 | | private static partial class Log |
| | | 81 | | { |
| | | 82 | | [LoggerMessage(Level = LogLevel.Debug, Message = "Wrote {ProxyType} proxy ({SizeBytes} B) for {RelativePath}.")] |
| | | 83 | | public static partial void ProxyWritten(ILogger logger, ProxyType proxyType, long sizeBytes, string relativePath |
| | | 84 | | |
| | | 85 | | [LoggerMessage(Level = LogLevel.Warning, Message = "Could not delete temporary proxy file {Path}.")] |
| | | 86 | | public static partial void DiscardFailed(ILogger logger, string path, Exception ex); |
| | | 87 | | } |
| | | 88 | | } |