diff --git a/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs b/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs index 06b68b114..6377d62f7 100644 --- a/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs +++ b/services/Archival/FilterLists.Archival.Application/Commands/ArchiveList.cs @@ -2,10 +2,10 @@ using System.Collections.Generic; using System.IO; using System.Linq; -using System.Net.Http; using System.Threading; using System.Threading.Tasks; using FilterLists.Archival.Application.Models; +using FilterLists.Archival.Infrastructure.Clients; using FilterLists.Archival.Infrastructure.Persistence; using FilterLists.SharedKernel.Apis.Clients; using FilterLists.SharedKernel.Apis.Contracts.Directory; @@ -29,19 +29,19 @@ public Command(int listId) public class Handler : IRequestHandler { private readonly IFileArchiver _archiver; + private readonly IFileClient _client; private readonly IDirectoryApi _directory; - private readonly HttpClient _httpClient; private readonly ILogger _logger; public Handler( IFileArchiver archiver, + IFileClient fileClient, IDirectoryApi directory, - IHttpClientFactory httpClientFactory, ILogger logger) { _archiver = archiver; + _client = fileClient; _directory = directory; - _httpClient = httpClientFactory.CreateClient(); _logger = logger; } @@ -86,34 +86,20 @@ private async Task DownloadSegments( IReadOnlyCollection segments, CancellationToken cancellationToken) { - var responses = new List(); - try + var downloads = new List>(); + foreach (var segment in segments) { - var readTasks = new List>(); - foreach (var segment in segments) - { - // TODO: add Polly for resiliency - var response = await _httpClient.GetAsync(segment.Url, HttpCompletionOption.ResponseHeadersRead, cancellationToken); - responses.Add(response); - response.EnsureSuccessStatusCode(); - readTasks.Add(response.Content.ReadAsStreamAsync()); - } - - var streams = await Task.WhenAll(readTasks); - - // TODO: prefix fileName with listId - // TODO: add ".txt" for sources with no extension - var fileName = Uri.UnescapeDataString(segments.First().Url.Segments.Last()); - var target = new FileInfo(fileName); - await _archiver.ArchiveFileAsync(new FileToArchive(target, streams), cancellationToken); - } - finally - { - foreach (var response in responses) - { - response.Dispose(); - } + downloads.Add(_client.DownloadFileAsync(segment.Url, cancellationToken)); } + + var streams = await Task.WhenAll(downloads); + + // TODO: prefix fileName with listId + // TODO: add ".txt" for sources with no extension + var fileName = Uri.UnescapeDataString(segments.First().Url.Segments.Last()); + var target = new FileInfo(fileName); + var file = new FileToArchive(target, streams); + await _archiver.ArchiveFileAsync(file, cancellationToken); } } } diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Clients/ConfigurationExtensions.cs b/services/Archival/FilterLists.Archival.Infrastructure/Clients/ConfigurationExtensions.cs new file mode 100644 index 000000000..f1119b2a2 --- /dev/null +++ b/services/Archival/FilterLists.Archival.Infrastructure/Clients/ConfigurationExtensions.cs @@ -0,0 +1,18 @@ +using System; +using Microsoft.Extensions.DependencyInjection; +using Polly; + +namespace FilterLists.Archival.Infrastructure.Clients +{ + internal static class ConfigurationExtensions + { + public static void AddClients(this IServiceCollection services) + { + services.AddHttpClient() + .AddTransientHttpErrorPolicy(b => b.WaitAndRetryAsync(new[] + { + TimeSpan.FromSeconds(1), TimeSpan.FromSeconds(5), TimeSpan.FromSeconds(10) + })); + } + } +} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Clients/FileClient.cs b/services/Archival/FilterLists.Archival.Infrastructure/Clients/FileClient.cs new file mode 100644 index 000000000..003c5c04f --- /dev/null +++ b/services/Archival/FilterLists.Archival.Infrastructure/Clients/FileClient.cs @@ -0,0 +1,38 @@ +using System; +using System.Collections.Generic; +using System.IO; +using System.Net.Http; +using System.Threading; +using System.Threading.Tasks; + +namespace FilterLists.Archival.Infrastructure.Clients +{ + internal sealed class FileClient : IFileClient + { + private readonly HttpClient _httpClient; + private readonly ICollection _httpResponseMessages = new List(); + + public FileClient(HttpClient httpClient) + { + _httpClient = httpClient; + } + + public async Task DownloadFileAsync(Uri url, CancellationToken cancellationToken) + { + var response = await _httpClient.GetAsync(url, HttpCompletionOption.ResponseHeadersRead, cancellationToken); + _httpResponseMessages.Add(response); + response.EnsureSuccessStatusCode(); + return await response.Content.ReadAsStreamAsync(); + } + + public void Dispose() + { + foreach (var message in _httpResponseMessages) + { + message.Dispose(); + } + + _httpClient.Dispose(); + } + } +} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/Clients/IFileClient.cs b/services/Archival/FilterLists.Archival.Infrastructure/Clients/IFileClient.cs new file mode 100644 index 000000000..b4d6a1ceb --- /dev/null +++ b/services/Archival/FilterLists.Archival.Infrastructure/Clients/IFileClient.cs @@ -0,0 +1,12 @@ +using System; +using System.IO; +using System.Threading; +using System.Threading.Tasks; + +namespace FilterLists.Archival.Infrastructure.Clients +{ + public interface IFileClient : IDisposable + { + Task DownloadFileAsync(Uri url, CancellationToken cancellationToken); + } +} diff --git a/services/Archival/FilterLists.Archival.Infrastructure/ConfigurationExtensions.cs b/services/Archival/FilterLists.Archival.Infrastructure/ConfigurationExtensions.cs index a20ec9d45..d801483c0 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/ConfigurationExtensions.cs +++ b/services/Archival/FilterLists.Archival.Infrastructure/ConfigurationExtensions.cs @@ -1,4 +1,5 @@ using System; +using FilterLists.Archival.Infrastructure.Clients; using FilterLists.Archival.Infrastructure.Persistence; using FilterLists.Archival.Infrastructure.Scheduling; using FilterLists.SharedKernel.Apis.Clients; @@ -24,6 +25,7 @@ public static void AddInfrastructureServices(this IServiceCollection services, I services.AddSharedKernelLogging(configuration); services.AddSchedulingServices(configuration); services.AddApiClients(configuration); + services.AddClients(); services.AddPersistenceServices(configuration); } diff --git a/services/Archival/FilterLists.Archival.Infrastructure/FilterLists.Archival.Infrastructure.csproj b/services/Archival/FilterLists.Archival.Infrastructure/FilterLists.Archival.Infrastructure.csproj index 31cf441a1..ad098ebbe 100644 --- a/services/Archival/FilterLists.Archival.Infrastructure/FilterLists.Archival.Infrastructure.csproj +++ b/services/Archival/FilterLists.Archival.Infrastructure/FilterLists.Archival.Infrastructure.csproj @@ -25,6 +25,7 @@ all runtime; build; native; contentfiles; analyzers; buildtransitive +