From e95892af1bc07af60642ee60c21e1617f80c5753 Mon Sep 17 00:00:00 2001 From: "Collin M. Barrett" Date: Fri, 18 Sep 2020 16:57:28 -0500 Subject: [PATCH] =?UTF-8?q?feat(archival):=20=E2=9C=A8=E2=9A=A1=E2=99=BB?= =?UTF-8?q?=20support=20segments=20of=20different=20source=20types,=20use?= =?UTF-8?q?=20IAsyncEnumerable=20for=20Streams?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Commands/ArchiveList.cs | 35 +++++++-------- .../Models/FileToArchive.cs | 20 ++++++--- .../FileStreamConversionStrategyFactory.cs | 23 ++++++++++ .../FileWriteStrategyFactory.cs | 22 --------- .../IFileStreamConversionStrategy.cs | 10 +++++ .../FileWriteStrategies/IFileWriteStrategy.cs | 10 ----- .../Persistence/FileWriteStrategies/Txt.cs | 14 +++--- .../Persistence/GitFileArchiver.cs | 45 +++++++++++-------- .../Persistence/IFileArchiver.cs | 2 +- .../Persistence/IFileToArchive.cs | 10 +++-- .../IDirectoryApi.cs | 4 +- 11 files changed, 106 insertions(+), 89 deletions(-) create mode 100644 services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileStreamConversionStrategyFactory.cs delete mode 100644 services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileWriteStrategyFactory.cs create mode 100644 services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileStreamConversionStrategy.cs delete mode 100644 services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileWriteStrategy.cs diff --git a/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs b/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs index ddae8d05f..0b1cb0722 100644 --- a/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs +++ b/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.IO; using System.Linq; +using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; using FilterLists.Archival.Application.Models; @@ -53,7 +54,7 @@ public async Task Handle(Command request, CancellationToken cancellationTo var segmentUrls = await GetSegmentUrlsAsync(request.ListId, cancellationToken); if (segmentUrls.Count > 0) { - var file = await GetFileAsync(request.ListId, segmentUrls, cancellationToken); + var file = GetFileToArchive(request.ListId, segmentUrls, cancellationToken); await _archiver.ArchiveFileAsync(file, cancellationToken); _archiver.Commit(); @@ -83,30 +84,28 @@ private async Task> GetSegmentUrlsAsync( new List(); } - private async Task GetFileAsync( + private IFileToArchive GetFileToArchive( int listId, - IReadOnlyCollection segmentUrls, - CancellationToken cancellationToken) - { - var contentsAsync = GetContentsAsync(segmentUrls, cancellationToken); - var sourceFileName = Uri.UnescapeDataString(segmentUrls.First().Url.Segments.Last()); - var sourceExtension = Path.GetExtension(sourceFileName); - var target = new FileInfo($"{listId}.txt"); - return new FileToArchive(sourceExtension, await contentsAsync, target); - } - - // TODO: IAsyncEnumerable? - private async Task> GetContentsAsync( IEnumerable segmentUrls, CancellationToken cancellationToken) { - var downloads = new List>(); + var segmentsAsync = GetSegmentsAsync(segmentUrls, cancellationToken); + var target = new FileInfo($"{listId}.txt"); + return new FileToArchive(segmentsAsync, target); + } + + private async IAsyncEnumerable GetSegmentsAsync( + IEnumerable segmentUrls, + [EnumeratorCancellation] CancellationToken cancellationToken) + { foreach (var segment in segmentUrls) { - downloads.Add(_client.DownloadFileAsync(segment.Url, cancellationToken)); + var sourceFileName = Uri.UnescapeDataString(segment.Url.Segments.Last()); + var sourceExtension = Path.GetExtension(sourceFileName); + yield return new FileToArchiveSegment( + sourceExtension, + await _client.DownloadFileAsync(segment.Url, cancellationToken)); } - - return await Task.WhenAll(downloads); } } } diff --git a/services/Archival/FilterLists.Archival.Application/Models/FileToArchive.cs b/services/Archival/FilterLists.Archival.Application/Models/FileToArchive.cs index 26949b566..dc217490b 100644 --- a/services/Archival/FilterLists.Archival.Application/Models/FileToArchive.cs +++ b/services/Archival/FilterLists.Archival.Application/Models/FileToArchive.cs @@ -6,15 +6,25 @@ namespace FilterLists.Archival.Application.Models { internal class FileToArchive : IFileToArchive { - public FileToArchive(string sourceExtension, IEnumerable contents, FileInfo target) + public FileToArchive(IAsyncEnumerable segments, FileInfo target) { - SourceExtension = sourceExtension; - Contents = contents; + Segments = segments; Target = target; } - public string SourceExtension { get; } - public IEnumerable Contents { get; } + public IAsyncEnumerable Segments { get; } public FileInfo Target { get; } } + + internal class FileToArchiveSegment : IFileToArchiveSegment + { + public FileToArchiveSegment(string sourceExtension, Stream contents) + { + SourceExtension = sourceExtension; + Contents = contents; + } + + public string SourceExtension { get; } + public Stream Contents { get; } + } } diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileStreamConversionStrategyFactory.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileStreamConversionStrategyFactory.cs new file mode 100644 index 000000000..303a8a1f7 --- /dev/null +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileStreamConversionStrategyFactory.cs @@ -0,0 +1,23 @@ +using System; +using System.Linq; +using System.Reflection; + +namespace FilterLists.Archival.Infrastructure.Persistence.FileWriteStrategies +{ + internal static class FileStreamConversionStrategyFactory + { + public static TFileStreamConversionStrategy? GetStrategy( + this IFileToArchiveSegment segment) + where TFileStreamConversionStrategy : class, IFileStreamConversionStrategy + { + var strategyType = Assembly.GetExecutingAssembly() + .GetTypes() + .FirstOrDefault(t => + typeof(TFileStreamConversionStrategy).IsAssignableFrom(t) && + string.Equals(t.Name, segment.SourceExtension.TrimStart('.'), StringComparison.OrdinalIgnoreCase)); + return strategyType is default(Type) + ? default + : (TFileStreamConversionStrategy)Activator.CreateInstance(strategyType); + } + } +} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileWriteStrategyFactory.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileWriteStrategyFactory.cs deleted file mode 100644 index f53655881..000000000 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/FileWriteStrategyFactory.cs +++ /dev/null @@ -1,22 +0,0 @@ -using System; -using System.Linq; -using System.Reflection; - -namespace FilterLists.Archival.Infrastructure.Persistence.FileWriteStrategies -{ - internal static class FileWriteStrategyFactory - { - public static TFileWriteStrategy? GetStrategy(this IFileToArchive file) - where TFileWriteStrategy : class, IFileWriteStrategy - { - var strategyType = Assembly.GetExecutingAssembly() - .GetTypes() - .FirstOrDefault(t => - typeof(TFileWriteStrategy).IsAssignableFrom(t) && - string.Equals(t.Name, file.SourceExtension.TrimStart('.'), StringComparison.OrdinalIgnoreCase)); - return strategyType is default(Type) - ? default - : (TFileWriteStrategy)Activator.CreateInstance(strategyType); - } - } -} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileStreamConversionStrategy.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileStreamConversionStrategy.cs new file mode 100644 index 000000000..5092ce368 --- /dev/null +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileStreamConversionStrategy.cs @@ -0,0 +1,10 @@ +using System.IO; +using System.Threading; + +namespace FilterLists.Archival.Infrastructure.Persistence.FileWriteStrategies +{ + internal interface IFileStreamConversionStrategy + { + Stream Convert(IFileToArchiveSegment segment, CancellationToken cancellationToken); + } +} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileWriteStrategy.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileWriteStrategy.cs deleted file mode 100644 index 5e6ec88e4..000000000 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/IFileWriteStrategy.cs +++ /dev/null @@ -1,10 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; - -namespace FilterLists.Archival.Infrastructure.Persistence.FileWriteStrategies -{ - internal interface IFileWriteStrategy - { - Task WriteAsync(IFileToArchive file, CancellationToken cancellationToken = default); - } -} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/Txt.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/Txt.cs index 840a0cbea..336bdd7f0 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/Txt.cs +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/FileWriteStrategies/Txt.cs @@ -1,20 +1,16 @@ using System; +using System.IO; using System.Threading; -using System.Threading.Tasks; namespace FilterLists.Archival.Infrastructure.Persistence.FileWriteStrategies { - public class Txt : IFileWriteStrategy + public class Txt : IFileStreamConversionStrategy { - public async Task WriteAsync(IFileToArchive file, CancellationToken cancellationToken = default) + public Stream Convert(IFileToArchiveSegment segment, CancellationToken cancellationToken) { - _ = file ?? throw new ArgumentNullException(nameof(file)); + _ = segment ?? throw new ArgumentNullException(nameof(segment)); - await using var target = file.Target.OpenWrite(); - foreach (var source in file.Contents) - { - await source.CopyToAsync(target, cancellationToken); - } + return segment.Contents; } } } diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/GitFileArchiver.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/GitFileArchiver.cs index 22a22dd3a..a12be7955 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/GitFileArchiver.cs +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/GitFileArchiver.cs @@ -31,29 +31,36 @@ public GitFileArchiver( public async Task ArchiveFileAsync( IFileToArchive file, - CancellationToken cancellationToken = default) + CancellationToken cancellationToken) { - var strategy = file.GetStrategy(); - if (strategy is default(IFileWriteStrategy)) + var textStreams = new List(); + await foreach (var segment in file.Segments.WithCancellation(cancellationToken)) { - _logger.LogWarning( - "No write strategy found for {FileName}. Skipping", - file.Target.Name); + var strategy = segment.GetStrategy(); + if (strategy is default(IFileStreamConversionStrategy)) + { + _logger.LogWarning( + "No file stream conversion strategy found for extension {Extension} in target {Target}. Skipping file", + segment.SourceExtension, + file.Target); + break; + } + + textStreams.Add(strategy.Convert(segment, cancellationToken)); } - else + + _logger.LogDebug("Writing {FileName}", file.Target.Name); + + _writtenFiles.Add(file.Target); + + // TODO: write to _options.RepositoryPath + await using var target = file.Target.OpenWrite(); + foreach (var textStream in textStreams) { - _logger.LogDebug( - "Writing {FileName} with strategy {FileWriteStrategy}", - file.Target.Name, - strategy.GetType().Name); - - _writtenFiles.Add(file.Target); - - // TODO: write to _options.RepositoryPath - await strategy.WriteAsync(file, cancellationToken); - - _logger.LogDebug("Finished writing {FileName}", file.Target.Name); + await textStream.CopyToAsync(target, cancellationToken); } + + _logger.LogDebug("Finished writing {FileName}", file.Target.Name); } public void Commit() @@ -61,7 +68,7 @@ public void Commit() var fileNames = _writtenFiles.Select(f => f.Name).ToList(); Commands.Stage(_repo, fileNames); var signature = new Signature(_options.UserName, _options.UserEmail, DateTime.UtcNow); - var message = $"feat(archives): archive {fileNames.Count} files{Environment.NewLine}{string.Join(Environment.NewLine, fileNames)}"; + var message = $"feat(archives): archive {fileNames.Count} files{Environment.NewLine}{string.Join(Environment.NewLine, fileNames)}"; _repo.Commit(message, signature, signature); _logger.LogDebug("Committed {@FileNames}", fileNames); diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileArchiver.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileArchiver.cs index 42a8fa844..1950c109b 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileArchiver.cs +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileArchiver.cs @@ -8,6 +8,6 @@ public interface IFileArchiver : IUnitOfWork { Task ArchiveFileAsync( IFileToArchive file, - CancellationToken cancellationToken = default); + CancellationToken cancellationToken); } } diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileToArchive.cs b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileToArchive.cs index 0776418b2..41ee26ea2 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileToArchive.cs +++ b/services/Archival/FilterLists.Archival.Infrastructure/Persistence/IFileToArchive.cs @@ -5,9 +5,13 @@ namespace FilterLists.Archival.Infrastructure.Persistence { public interface IFileToArchive { - // for now, assume all source segments have same extension - string SourceExtension { get; } - IEnumerable Contents { get; } + IAsyncEnumerable Segments { get; } FileInfo Target { get; } } + + public interface IFileToArchiveSegment + { + string SourceExtension { get; } + Stream Contents { get; } + } } diff --git a/services/SharedKernel/FilterLists.SharedKernel.Apis.Clients/IDirectoryApi.cs b/services/SharedKernel/FilterLists.SharedKernel.Apis.Clients/IDirectoryApi.cs index efb731560..9fcd94c6c 100644 --- a/services/SharedKernel/FilterLists.SharedKernel.Apis.Clients/IDirectoryApi.cs +++ b/services/SharedKernel/FilterLists.SharedKernel.Apis.Clients/IDirectoryApi.cs @@ -9,9 +9,9 @@ namespace FilterLists.SharedKernel.Apis.Clients public interface IDirectoryApi { [Get("/lists")] - Task> GetListsAsync(CancellationToken cancellationToken = default); + Task> GetListsAsync(CancellationToken cancellationToken); [Get("/lists/{id}")] - Task GetListDetailsAsync(int id, CancellationToken cancellationToken = default); + Task GetListDetailsAsync(int id, CancellationToken cancellationToken); } }