diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props index 30d13cbed0..cfdce8ef0a 100644 --- a/src/Directory.Packages.props +++ b/src/Directory.Packages.props @@ -6,9 +6,11 @@ + + @@ -76,6 +78,8 @@ + + diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistence.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistence.cs index dabe754e30..9e807fc358 100644 --- a/src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistence.cs +++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistence.cs @@ -12,7 +12,7 @@ public void AddPersistence(IServiceCollection services) { RegisterSettings(services); ConfigureDbContext(services); - RegisterDataStores(services); + RegisterDataStores(services, settings); services.AddSingleton(); } @@ -23,6 +23,7 @@ public void AddInstaller(IServiceCollection services) ConfigureDbContext(services); services.AddScoped(); + RegisterBodyStorageInstaller(services, settings); } void RegisterSettings(IServiceCollection services) diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/persistence.manifest b/src/ServiceControl.Persistence.EFCore.PostgreSql/persistence.manifest index 925712ab78..8e667275c6 100644 --- a/src/ServiceControl.Persistence.EFCore.PostgreSql/persistence.manifest +++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/persistence.manifest @@ -13,10 +13,58 @@ "Name": "ServiceControl/Database/CommandTimeout", "Mandatory": false }, + { + "Name": "ServiceControl/MessageBody/StorageType", + "Mandatory": true + }, { "Name": "ServiceControl/MessageBody/StoragePath", "Mandatory": false }, + { + "Name": "ServiceControl/MessageBody/Azure/ConnectionString", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ServiceUri", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ManagedIdentityClientId", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/AuthorityHost", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ContainerName", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/BucketName", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/KeyPrefix", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/Region", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/ServiceUrl", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/AccessKeyId", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/SecretAccessKey", + "Mandatory": false + }, { "Name": "ServiceControl/MessageBody/MinCompressionSize", "Mandatory": false diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistence.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistence.cs index 5e9336e747..3c6d1384b3 100644 --- a/src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistence.cs +++ b/src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistence.cs @@ -12,7 +12,7 @@ public void AddPersistence(IServiceCollection services) { RegisterSettings(services); ConfigureDbContext(services); - RegisterDataStores(services); + RegisterDataStores(services, settings); services.AddSingleton(); } @@ -23,6 +23,7 @@ public void AddInstaller(IServiceCollection services) ConfigureDbContext(services); services.AddScoped(); + RegisterBodyStorageInstaller(services, settings); } void RegisterSettings(IServiceCollection services) diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/persistence.manifest b/src/ServiceControl.Persistence.EFCore.SqlServer/persistence.manifest index aa1b57b88f..51de2f2ee8 100644 --- a/src/ServiceControl.Persistence.EFCore.SqlServer/persistence.manifest +++ b/src/ServiceControl.Persistence.EFCore.SqlServer/persistence.manifest @@ -13,10 +13,58 @@ "Name": "ServiceControl/Database/CommandTimeout", "Mandatory": false }, + { + "Name": "ServiceControl/MessageBody/StorageType", + "Mandatory": true + }, { "Name": "ServiceControl/MessageBody/StoragePath", "Mandatory": false }, + { + "Name": "ServiceControl/MessageBody/Azure/ConnectionString", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ServiceUri", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ManagedIdentityClientId", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/AuthorityHost", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/Azure/ContainerName", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/BucketName", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/KeyPrefix", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/Region", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/ServiceUrl", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/AccessKeyId", + "Mandatory": false + }, + { + "Name": "ServiceControl/MessageBody/S3/SecretAccessKey", + "Mandatory": false + }, { "Name": "ServiceControl/MessageBody/MinCompressionSize", "Mandatory": false diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs index c4d67e60dc..c942102e14 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs @@ -13,7 +13,7 @@ namespace ServiceControl.Persistence.EFCore.Abstractions; public abstract class BasePersistence { - protected static void RegisterDataStores(IServiceCollection services) + protected static void RegisterDataStores(IServiceCollection services, EFPersisterSettings settings) { services.AddSingleton(TimeProvider.System); services.AddSingleton(); @@ -51,6 +51,43 @@ protected static void RegisterDataStores(IServiceCollection services) services.AddSingleton(); - services.AddSingleton(); + RegisterBodyStorage(services, settings); + } + + static void RegisterBodyStorage(IServiceCollection services, EFPersisterSettings settings) + { + switch (settings.BodyStorageType) + { + case BodyStorageType.FileSystem: + services.AddSingleton(); + break; + case BodyStorageType.AzureBlob: + services.AddSingleton(); + break; + case BodyStorageType.S3: + services.AddSingleton(); + break; + default: + throw new ArgumentOutOfRangeException(nameof(settings), settings.BodyStorageType, "Unknown body storage type."); + } + } + + // Only stores needing setup-time provisioning register an installer; SetupCommand skips when none is. + protected static void RegisterBodyStorageInstaller(IServiceCollection services, EFPersisterSettings settings) + { + switch (settings.BodyStorageType) + { + case BodyStorageType.FileSystem: + services.AddScoped(); + break; + case BodyStorageType.AzureBlob: + services.AddScoped(); + break; + case BodyStorageType.S3: + services.AddScoped(); + break; + default: + throw new ArgumentOutOfRangeException(nameof(settings), settings.BodyStorageType, "Unknown body storage type."); + } } } diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/BodyStorageType.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/BodyStorageType.cs new file mode 100644 index 0000000000..9485cab036 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/BodyStorageType.cs @@ -0,0 +1,8 @@ +namespace ServiceControl.Persistence.EFCore.Abstractions; + +public enum BodyStorageType +{ + FileSystem, + AzureBlob, + S3 +} diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs index fba0f10e3d..51252d4b60 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs @@ -8,7 +8,19 @@ public abstract class EFPersistenceConfigurationBase : IPersistenceConfiguration { const string ConnectionStringKey = "Database/ConnectionString"; const string CommandTimeoutKey = "Database/CommandTimeout"; + const string BodyStorageTypeKey = "MessageBody/StorageType"; const string MessageBodyStoragePathKey = "MessageBody/StoragePath"; + const string AzureConnectionStringKey = "MessageBody/Azure/ConnectionString"; + const string AzureServiceUriKey = "MessageBody/Azure/ServiceUri"; + const string AzureManagedIdentityClientIdKey = "MessageBody/Azure/ManagedIdentityClientId"; + const string AzureAuthorityHostKey = "MessageBody/Azure/AuthorityHost"; + const string AzureContainerNameKey = "MessageBody/Azure/ContainerName"; + const string S3BucketNameKey = "MessageBody/S3/BucketName"; + const string S3KeyPrefixKey = "MessageBody/S3/KeyPrefix"; + const string S3RegionKey = "MessageBody/S3/Region"; + const string S3ServiceUrlKey = "MessageBody/S3/ServiceUrl"; + const string S3AccessKeyIdKey = "MessageBody/S3/AccessKeyId"; + const string S3SecretAccessKeyKey = "MessageBody/S3/SecretAccessKey"; const string MinBodySizeForCompressionKey = "MessageBody/MinCompressionSize"; const string MaxBodySizeToStoreKey = "MaxBodySizeToStore"; const string ErrorRetentionPeriodKey = "ErrorRetentionPeriod"; @@ -19,12 +31,13 @@ public PersistenceSettings CreateSettings(SettingsRootNamespace settingsRootName var settings = CreateSettings(GetRequiredSetting(settingsRootNamespace, ConnectionStringKey)); settings.CommandTimeout = SettingsReader.Read(settingsRootNamespace, CommandTimeoutKey, EFPersisterSettings.DefaultCommandTimeout); - settings.MessageBodyStoragePath = SettingsReader.Read(settingsRootNamespace, MessageBodyStoragePathKey); settings.MinBodySizeForCompression = SettingsReader.Read(settingsRootNamespace, MinBodySizeForCompressionKey, EFPersisterSettings.DefaultMinBodySizeForCompression); settings.MaxBodySizeToStore = ReadMaxBodySizeToStore(settingsRootNamespace); settings.ErrorRetentionPeriod = GetRequiredSetting(settingsRootNamespace, ErrorRetentionPeriodKey); settings.EnableFullTextSearchOnBodies = SettingsReader.Read(settingsRootNamespace, EnableFullTextSearchOnBodiesKey, true); + ConfigureBodyStorage(settings, settingsRootNamespace); + return settings; } @@ -32,6 +45,101 @@ public PersistenceSettings CreateSettings(SettingsRootNamespace settingsRootName protected abstract EFPersisterSettings CreateSettings(string connectionString); + static void ConfigureBodyStorage(EFPersisterSettings settings, SettingsRootNamespace settingsRootNamespace) + { + settings.BodyStorageType = ReadBodyStorageType(settingsRootNamespace); + + if (settings.BodyStorageType == BodyStorageType.FileSystem) + { + settings.MessageBodyStoragePath = GetRequiredSetting(settingsRootNamespace, MessageBodyStoragePathKey); + } + else if (settings.BodyStorageType == BodyStorageType.AzureBlob) + { + ConfigureAzureBlob(settings, settingsRootNamespace); + } + else if (settings.BodyStorageType == BodyStorageType.S3) + { + ConfigureS3(settings, settingsRootNamespace); + } + } + + static void ConfigureS3(EFPersisterSettings settings, SettingsRootNamespace settingsRootNamespace) + { + settings.S3BucketName = GetRequiredSetting(settingsRootNamespace, S3BucketNameKey); + settings.S3KeyPrefix = SettingsReader.Read(settingsRootNamespace, S3KeyPrefixKey, settings.S3KeyPrefix); + settings.S3Region = SettingsReader.Read(settingsRootNamespace, S3RegionKey); + settings.S3ServiceUrl = SettingsReader.Read(settingsRootNamespace, S3ServiceUrlKey); + + var accessKeyId = SettingsReader.Read(settingsRootNamespace, S3AccessKeyIdKey); + var secretAccessKey = SettingsReader.Read(settingsRootNamespace, S3SecretAccessKeyKey); + var hasAccessKeyId = !string.IsNullOrWhiteSpace(accessKeyId); + var hasSecretAccessKey = !string.IsNullOrWhiteSpace(secretAccessKey); + + if (hasAccessKeyId != hasSecretAccessKey) + { + throw new Exception($"S3 body storage requires both {S3AccessKeyIdKey} and {S3SecretAccessKeyKey} together, or neither to use the default AWS credential chain (IAM role)."); + } + + settings.S3AccessKeyId = hasAccessKeyId ? accessKeyId : null; + settings.S3SecretAccessKey = hasSecretAccessKey ? secretAccessKey : null; + } + + static void ConfigureAzureBlob(EFPersisterSettings settings, SettingsRootNamespace settingsRootNamespace) + { + var connectionString = SettingsReader.Read(settingsRootNamespace, AzureConnectionStringKey); + var serviceUri = SettingsReader.Read(settingsRootNamespace, AzureServiceUriKey); + + var hasConnectionString = !string.IsNullOrWhiteSpace(connectionString); + var hasServiceUri = !string.IsNullOrWhiteSpace(serviceUri); + + // A connection string carries shared-key/SAS auth; a service URI uses managed identity. They + // are mutually exclusive, and exactly one must be provided. + if (hasConnectionString == hasServiceUri) + { + throw new Exception($"Azure Blob body storage requires exactly one of {AzureConnectionStringKey} (shared key / SAS) or {AzureServiceUriKey} (managed identity). {(hasConnectionString ? "Both were set." : "Neither was set.")}"); + } + + settings.AzureBlobConnectionString = hasConnectionString ? connectionString : null; + settings.AzureBlobServiceUri = hasServiceUri ? serviceUri : null; + settings.AzureBlobManagedIdentityClientId = SettingsReader.Read(settingsRootNamespace, AzureManagedIdentityClientIdKey); + settings.AzureBlobAuthorityHost = ReadAzureAuthorityHost(settingsRootNamespace); + settings.AzureBlobContainerName = SettingsReader.Read(settingsRootNamespace, AzureContainerNameKey, settings.AzureBlobContainerName); + } + + static string? ReadAzureAuthorityHost(SettingsRootNamespace settingsRootNamespace) + { + var raw = SettingsReader.Read(settingsRootNamespace, AzureAuthorityHostKey); + + if (string.IsNullOrWhiteSpace(raw)) + { + return null; + } + + if (!Uri.TryCreate(raw, UriKind.Absolute, out _)) + { + throw new Exception($"Setting {AzureAuthorityHostKey} value '{raw}' is not a valid absolute URI."); + } + + return raw; + } + + static BodyStorageType ReadBodyStorageType(SettingsRootNamespace settingsRootNamespace) + { + var raw = SettingsReader.Read(settingsRootNamespace, BodyStorageTypeKey); + + if (string.IsNullOrWhiteSpace(raw)) + { + throw new Exception($"Setting {BodyStorageTypeKey} is required. Valid values: {string.Join(", ", Enum.GetNames())}."); + } + + if (!Enum.TryParse(raw, ignoreCase: true, out var value) || !Enum.IsDefined(value)) + { + throw new Exception($"Setting {BodyStorageTypeKey} value '{raw}' is not valid. Valid values: {string.Join(", ", Enum.GetNames())}."); + } + + return value; + } + static int ReadMaxBodySizeToStore(SettingsRootNamespace settingsRootNamespace) { var maxBodySizeToStore = SettingsReader.Read(settingsRootNamespace, MaxBodySizeToStoreKey, EFPersisterSettings.DefaultMaxBodySizeToStore); diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersisterSettings.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersisterSettings.cs index 3b919453c4..6401d7267f 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersisterSettings.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersisterSettings.cs @@ -11,7 +11,19 @@ public abstract class EFPersisterSettings : PersistenceSettings public required string ConnectionString { get; set; } public int CommandTimeout { get; set; } = DefaultCommandTimeout; public TimeSpan ErrorRetentionPeriod { get; set; } + public BodyStorageType BodyStorageType { get; set; } public string? MessageBodyStoragePath { get; set; } + public string? AzureBlobConnectionString { get; set; } + public string? AzureBlobServiceUri { get; set; } + public string? AzureBlobManagedIdentityClientId { get; set; } + public string? AzureBlobAuthorityHost { get; set; } + public string AzureBlobContainerName { get; set; } = "error-bodies"; + public string? S3BucketName { get; set; } + public string S3KeyPrefix { get; set; } = "error-bodies/"; + public string? S3Region { get; set; } + public string? S3ServiceUrl { get; set; } + public string? S3AccessKeyId { get; set; } + public string? S3SecretAccessKey { get; set; } public int MinBodySizeForCompression { get; set; } = DefaultMinBodySizeForCompression; public int MaxBodySizeToStore { get; set; } = DefaultMaxBodySizeToStore; public int MaxRetryCount { get; set; } = 5; diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStorageInstaller.cs b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStorageInstaller.cs new file mode 100644 index 0000000000..e4050494db --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStorageInstaller.cs @@ -0,0 +1,10 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using ServiceControl.Persistence; +using ServiceControl.Persistence.EFCore.Abstractions; + +public class AzureBlobBodyStorageInstaller(EFPersisterSettings settings) : IBodyStorageInstaller +{ + public Task Provision(CancellationToken cancellationToken = default) => + AzureBlobClientFactory.CreateContainerClient(settings).CreateIfNotExistsAsync(cancellationToken: cancellationToken); +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStoragePersistence.cs b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStoragePersistence.cs new file mode 100644 index 0000000000..8e6aeed70a --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobBodyStoragePersistence.cs @@ -0,0 +1,98 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using System.IO.Compression; +using System.Net; +using Azure; +using Azure.Storage.Blobs; +using Azure.Storage.Blobs.Models; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Infrastructure; + +public class AzureBlobBodyStoragePersistence : IBodyStoragePersistence +{ + const string FormatVersion = "1"; + + readonly BlobContainerClient container; + readonly int minBodySizeForCompression; + + public AzureBlobBodyStoragePersistence(EFPersisterSettings settings) + { + container = AzureBlobClientFactory.CreateContainerClient(settings); + minBodySizeForCompression = settings.MinBodySizeForCompression; + } + + public async Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) + { + var blob = container.GetBlobClient(bodyId); + + var compressed = body.Length >= minBodySizeForCompression ? BodyCompression.TryCompress(body) : null; + var data = compressed is null ? BinaryData.FromBytes(body) : BinaryData.FromBytes(compressed); + + var options = new BlobUploadOptions + { + // Bodies are immutable, so create the blob only if it does not already exist. + Conditions = new BlobRequestConditions { IfNoneMatch = ETag.All }, + Metadata = new Dictionary + { + ["FormatVersion"] = FormatVersion, + ["ContentType"] = Uri.EscapeDataString(contentType), + ["BodySize"] = body.Length.ToString(), + ["IsCompressed"] = (compressed is not null).ToString() + } + }; + + try + { + await blob.UploadAsync(data, options, cancellationToken); + } + catch (RequestFailedException ex) when (ex.Status is (int)HttpStatusCode.Conflict or (int)HttpStatusCode.PreconditionFailed) + { + // A blob for this immutable body already exists. + } + } + + public async Task ReadBody(string bodyId, CancellationToken cancellationToken = default) + { + var blob = container.GetBlobClient(bodyId); + + try + { + var content = (await blob.DownloadStreamingAsync(cancellationToken: cancellationToken)).Value; + try + { + var metadata = content.Details.Metadata; + + if (metadata.TryGetValue("FormatVersion", out var version) && version != FormatVersion) + { + throw new InvalidOperationException($"Unsupported blob format version {version} for {bodyId}."); + } + + var contentType = metadata.TryGetValue("ContentType", out var ct) ? Uri.UnescapeDataString(ct) : "application/octet-stream"; + var bodySize = metadata.TryGetValue("BodySize", out var sizeText) && int.TryParse(sizeText, out var size) ? size : 0; + var isCompressed = metadata.TryGetValue("IsCompressed", out var compressedText) && bool.TryParse(compressedText, out var compressed) && compressed; + Stream stream = isCompressed + ? new ExpectedLengthStream(new BrotliStream(content.Content, CompressionMode.Decompress), bodySize) + : content.Content; + + return new MessageBodyFileResult + { + Stream = stream, + ContentType = contentType, + BodySize = bodySize + }; + } + catch + { + content.Dispose(); + throw; + } + } + catch (RequestFailedException ex) when (ex.Status == (int)HttpStatusCode.NotFound) + { + return null; + } + } + + public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) => + container.GetBlobClient(bodyId).DeleteIfExistsAsync(cancellationToken: cancellationToken); +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobClientFactory.cs b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobClientFactory.cs new file mode 100644 index 0000000000..33666bcb1d --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/AzureBlobClientFactory.cs @@ -0,0 +1,37 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using Azure.Core; +using Azure.Identity; +using Azure.Storage.Blobs; +using ServiceControl.Persistence.EFCore.Abstractions; + +static class AzureBlobClientFactory +{ + public static BlobContainerClient CreateContainerClient(EFPersisterSettings settings) + { + var serviceClient = settings.AzureBlobConnectionString is { Length: > 0 } connectionString + ? new BlobServiceClient(connectionString) + : new BlobServiceClient(new Uri(settings.AzureBlobServiceUri!), CreateCredential(settings)); + + return serviceClient.GetBlobContainerClient(settings.AzureBlobContainerName); + } + + static TokenCredential CreateCredential(EFPersisterSettings settings) + { + var options = new DefaultAzureCredentialOptions(); + + // Steers the login endpoint for sovereign clouds; when unset the SDK honours the + // AZURE_AUTHORITY_HOST environment variable. + if (settings.AzureBlobAuthorityHost is { Length: > 0 } authorityHost) + { + options.AuthorityHost = new Uri(authorityHost); + } + + if (settings.AzureBlobManagedIdentityClientId is { Length: > 0 } clientId) + { + options.ManagedIdentityClientId = clientId; + } + + return new DefaultAzureCredential(options); + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/BodyCompression.cs b/src/ServiceControl.Persistence.EFCore/Implementation/BodyCompression.cs new file mode 100644 index 0000000000..49dbaf6277 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/BodyCompression.cs @@ -0,0 +1,36 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using System.Buffers; +using System.IO.Compression; + +static class BodyCompression +{ + // Returns the Brotli-compressed bytes, or null when compression fails and the caller should + // store the body uncompressed. + public static byte[]? TryCompress(ReadOnlyMemory body) + { + var buffer = ArrayPool.Shared.Rent(BrotliEncoder.GetMaxCompressedLength(body.Length)); + try + { + return BrotliEncoder.TryCompress(body.Span, buffer, out var written, quality: 1, window: 22) + ? buffer.AsSpan(0, written).ToArray() + : null; + } + finally + { + ArrayPool.Shared.Return(buffer); + } + } + + public static byte[] Decompress(ReadOnlySpan compressed, int originalSize) + { + var decompressed = new byte[originalSize]; + + if (!BrotliDecoder.TryDecompress(compressed, decompressed, out var written) || written != originalSize) + { + throw new InvalidOperationException("Failed to decompress the message body."); + } + + return decompressed; + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage.cs b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage.cs index e1caa893a7..6fbe29a8cd 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage.cs @@ -1,9 +1,97 @@ namespace ServiceControl.Persistence.EFCore.Implementation; +using System.Linq.Expressions; +using System.Text; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; using ServiceControl.Operations.BodyStorage; +using ServiceControl.Persistence.EFCore.DbContexts; +using ServiceControl.Persistence.EFCore.Entities; +using ServiceControl.Persistence.EFCore.Infrastructure; -public class BodyStorage : IBodyStorage +// A body is stored inline in BodyText (small text) or in external storage (binary, or large text +// where BodyText keeps only a search prefix). External storage is authoritative, so check it first. +// bodyId is usually a UniqueMessageId (a Guid) but may be a plain MessageId. +public class BodyStorage(IServiceScopeFactory scopeFactory, IBodyStoragePersistence storagePersistence) : DataStoreBase(scopeFactory), IBodyStorage { - public Task TryFetch(string bodyId) => - throw new NotImplementedException(); + public async Task TryFetch(string bodyId) + { + var row = await ExecuteWithDbContext(dbContext => ResolveBody(dbContext, bodyId)); + + if (row == null) + { + return null; // No such message: the API turns this into a 404. + } + + // Bodies are immutable per message, so the id is a stable ETag. + var uniqueMessageId = row.UniqueMessageId.ToString(); + + if (row.BodyStoredExternally) + { + var external = await storagePersistence.ReadBody(uniqueMessageId); + + return external == null + ? new MessageBodyStreamResult { HasResult = false } + : new MessageBodyStreamResult + { + HasResult = true, + Stream = external.Stream, + ContentType = external.ContentType, + BodySize = external.BodySize, + Etag = uniqueMessageId + }; + } + + if (row.BodyText != null) + { + var bytes = Encoding.UTF8.GetBytes(row.BodyText); + + return new MessageBodyStreamResult + { + HasResult = true, + Stream = new MemoryStream(bytes, writable: false), + ContentType = row.BodyContentType, + BodySize = bytes.Length, + Etag = uniqueMessageId + }; + } + + return new MessageBodyStreamResult { HasResult = false }; // Message exists but carries no body. + } + + static async Task ResolveBody(ServiceControlDbContext dbContext, string bodyId) + { + if (Guid.TryParse(bodyId, out var uniqueMessageId)) + { + var byUniqueId = await Query(dbContext, message => message.UniqueMessageId == uniqueMessageId); + if (byUniqueId != null) + { + return byUniqueId; + } + } + + return await Query(dbContext, message => message.MessageId == bodyId); + } + + static Task Query(ServiceControlDbContext dbContext, Expression> predicate) => + dbContext.FailedMessages + .AsNoTracking() + .Where(predicate) + .OrderBy(message => message.UniqueMessageId) + .Select(message => new BodyRow + { + UniqueMessageId = message.UniqueMessageId, + BodyText = message.BodyText, + BodyStoredExternally = message.BodyStoredExternally, + BodyContentType = message.BodyContentType + }) + .FirstOrDefaultAsync(); + + sealed class BodyRow + { + public Guid UniqueMessageId { get; init; } + public string? BodyText { get; init; } + public bool BodyStoredExternally { get; init; } + public string? BodyContentType { get; init; } + } } diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FakeBodyStoragePersistence.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FakeBodyStoragePersistence.cs deleted file mode 100644 index c4fe41a092..0000000000 --- a/src/ServiceControl.Persistence.EFCore/Implementation/FakeBodyStoragePersistence.cs +++ /dev/null @@ -1,10 +0,0 @@ -namespace ServiceControl.Persistence.EFCore.Implementation; - -using ServiceControl.Persistence.EFCore.Infrastructure; - -public class FakeBodyStoragePersistence : IBodyStoragePersistence -{ - public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) => throw new NotImplementedException(); - public Task ReadBody(string bodyId, CancellationToken cancellationToken = default) => throw new NotImplementedException(); - public Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) => throw new NotImplementedException(); -} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStorageInstaller.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStorageInstaller.cs new file mode 100644 index 0000000000..91fc120528 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStorageInstaller.cs @@ -0,0 +1,14 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using ServiceControl.Persistence; +using ServiceControl.Persistence.EFCore.Abstractions; + +public class FileSystemBodyStorageInstaller(EFPersisterSettings settings) : IBodyStorageInstaller +{ + public Task Provision(CancellationToken cancellationToken = default) + { + Directory.CreateDirectory(settings.MessageBodyStoragePath!); + + return Task.CompletedTask; + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStoragePersistence.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStoragePersistence.cs new file mode 100644 index 0000000000..9afe81b384 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/FileSystemBodyStoragePersistence.cs @@ -0,0 +1,146 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using System.IO.Compression; +using System.Text; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Infrastructure; + +// Bodies are immutable and keyed by bodyId alone, so a re-failure resolves to the same file and an +// existing one is left untouched. +public class FileSystemBodyStoragePersistence(EFPersisterSettings settings) : IBodyStoragePersistence +{ + const int FormatVersion = 1; + + string StoragePath => settings.MessageBodyStoragePath!; + + public async Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) + { + var filePath = GetBodyFilePath(bodyId); + + if (File.Exists(filePath)) + { + return; + } + + // A unique temp name lets concurrent writers of the same body race without clobbering. + var tempFilePath = $"{filePath}.{Guid.NewGuid():N}.tmp"; + + try + { + var fileStream = new FileStream(tempFilePath, FileMode.CreateNew, FileAccess.Write, FileShare.None, bufferSize: 4096, useAsync: true); + await using (fileStream.ConfigureAwait(false)) + { + var shouldCompress = body.Length >= settings.MinBodySizeForCompression; + + using (var writer = new BinaryWriter(fileStream, Encoding.UTF8, leaveOpen: true)) + { + writer.Write(FormatVersion); + writer.Write(contentType); + writer.Write(body.Length); + writer.Write(shouldCompress); + writer.Flush(); + } + + if (shouldCompress) + { + var brotliStream = new BrotliStream(fileStream, CompressionLevel.Fastest, leaveOpen: true); + await using (brotliStream.ConfigureAwait(false)) + { + await brotliStream.WriteAsync(body, cancellationToken).ConfigureAwait(false); + } + } + else + { + await fileStream.WriteAsync(body, cancellationToken).ConfigureAwait(false); + } + + await fileStream.FlushAsync(cancellationToken).ConfigureAwait(false); + } + + try + { + File.Move(tempFilePath, filePath, overwrite: false); + } + catch (IOException) when (File.Exists(filePath)) + { + // A concurrent writer already produced the (immutable) body; discard our copy. + TryDelete(tempFilePath); + } + } + catch + { + TryDelete(tempFilePath); + throw; + } + } + + public Task ReadBody(string bodyId, CancellationToken cancellationToken = default) + { + var filePath = GetBodyFilePath(bodyId); + + if (!File.Exists(filePath)) + { + return Task.FromResult(null); + } + + FileStream? fileStream = null; + try + { + fileStream = new FileStream(filePath, FileMode.Open, FileAccess.Read, FileShare.Read, bufferSize: 4096, useAsync: true); + + var reader = new BinaryReader(fileStream, Encoding.UTF8, leaveOpen: true); + + var formatVersion = reader.ReadInt32(); + if (formatVersion != FormatVersion) + { + throw new InvalidOperationException($"Unsupported body file format version {formatVersion} for {bodyId}."); + } + + var contentType = reader.ReadString(); + var bodySize = reader.ReadInt32(); + var isCompressed = reader.ReadBoolean(); + + // The returned stream owns fileStream and is disposed by the caller. + Stream bodyStream = isCompressed + ? new BrotliStream(fileStream, CompressionMode.Decompress, leaveOpen: false) + : fileStream; + + return Task.FromResult(new MessageBodyFileResult + { + Stream = bodyStream, + ContentType = contentType, + BodySize = bodySize + }); + } + catch (FileNotFoundException) + { + fileStream?.Dispose(); + return Task.FromResult(null); + } + catch + { + fileStream?.Dispose(); + throw; + } + } + + public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) + { + TryDelete(GetBodyFilePath(bodyId)); + return Task.CompletedTask; + } + + string GetBodyFilePath(string bodyId) => Path.Combine(StoragePath, $"{bodyId}.body"); + + static void TryDelete(string filePath) + { + try + { + File.Delete(filePath); + } + catch (DirectoryNotFoundException) + { + // Nothing to delete. + } + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStorageInstaller.cs b/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStorageInstaller.cs new file mode 100644 index 0000000000..1efb36f57e --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStorageInstaller.cs @@ -0,0 +1,19 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using Amazon.S3.Model; +using Amazon.S3.Util; +using ServiceControl.Persistence; +using ServiceControl.Persistence.EFCore.Abstractions; + +public class S3BodyStorageInstaller(EFPersisterSettings settings) : IBodyStorageInstaller +{ + public async Task Provision(CancellationToken cancellationToken = default) + { + using var client = S3ClientFactory.Create(settings); + + if (!await AmazonS3Util.DoesS3BucketExistV2Async(client, settings.S3BucketName)) + { + await client.PutBucketAsync(new PutBucketRequest { BucketName = settings.S3BucketName }, cancellationToken); + } + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStoragePersistence.cs b/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStoragePersistence.cs new file mode 100644 index 0000000000..10f6561108 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/S3BodyStoragePersistence.cs @@ -0,0 +1,113 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using System.IO.Compression; +using System.Net; +using Amazon.S3; +using Amazon.S3.Model; +using ServiceControl.Infrastructure; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Infrastructure; + +public class S3BodyStoragePersistence : IBodyStoragePersistence +{ + const string FormatVersion = "1"; + + readonly IAmazonS3 client; + readonly string bucketName; + readonly string keyPrefix; + readonly int minBodySizeForCompression; + + public S3BodyStoragePersistence(EFPersisterSettings settings) + { + client = S3ClientFactory.Create(settings); + bucketName = settings.S3BucketName!; + keyPrefix = settings.S3KeyPrefix; + minBodySizeForCompression = settings.MinBodySizeForCompression; + } + + public async Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) + { + var key = Key(bodyId); + + // Bodies are immutable, so an existing object is already correct. + if (await Exists(key, cancellationToken)) + { + return; + } + + var compressed = body.Length >= minBodySizeForCompression ? BodyCompression.TryCompress(body) : null; + + var payload = compressed is null ? body : compressed.AsMemory(); + var request = new PutObjectRequest + { + BucketName = bucketName, + Key = key, + InputStream = new ReadOnlyStream(payload), + ContentType = contentType + }; + request.Metadata.Add("format-version", FormatVersion); + request.Metadata.Add("body-size", body.Length.ToString()); + request.Metadata.Add("is-compressed", (compressed is not null).ToString()); + + await client.PutObjectAsync(request, cancellationToken); + } + + public async Task ReadBody(string bodyId, CancellationToken cancellationToken = default) + { + try + { + var response = await client.GetObjectAsync(bucketName, Key(bodyId), cancellationToken); + try + { + var metadata = response.Metadata; + + var version = metadata["format-version"]; + if (!string.IsNullOrEmpty(version) && version != FormatVersion) + { + throw new InvalidOperationException($"Unsupported object format version {version} for {bodyId}."); + } + + var bodySize = int.TryParse(metadata["body-size"], out var size) ? size : 0; + var isCompressed = bool.TryParse(metadata["is-compressed"], out var compressed) && compressed; + var contentType = response.Headers.ContentType ?? "application/octet-stream"; + Stream stream = isCompressed + ? new ExpectedLengthStream(new BrotliStream(response.ResponseStream, CompressionMode.Decompress), bodySize) + : response.ResponseStream; + + return new MessageBodyFileResult + { + Stream = new OwnedStream(stream, response), + ContentType = contentType, + BodySize = bodySize + }; + } + catch + { + response.Dispose(); + throw; + } + } + catch (AmazonS3Exception ex) when (ex.StatusCode == HttpStatusCode.NotFound) + { + return null; + } + } + + public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) => + client.DeleteObjectAsync(bucketName, Key(bodyId), cancellationToken); + + async Task Exists(string key, CancellationToken cancellationToken) + { + try + { + await client.GetObjectMetadataAsync(bucketName, key, cancellationToken); + return true; + } + catch (AmazonS3Exception ex) when (ex.StatusCode == HttpStatusCode.NotFound) + { + return false; + } + } + + string Key(string bodyId) => $"{keyPrefix}{bodyId}"; +} diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/S3ClientFactory.cs b/src/ServiceControl.Persistence.EFCore/Implementation/S3ClientFactory.cs new file mode 100644 index 0000000000..f1c6d24c6e --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Implementation/S3ClientFactory.cs @@ -0,0 +1,39 @@ +namespace ServiceControl.Persistence.EFCore.Implementation; + +using Amazon; +using Amazon.Runtime; +using Amazon.S3; +using ServiceControl.Persistence.EFCore.Abstractions; + +static class S3ClientFactory +{ + public static IAmazonS3 Create(EFPersisterSettings settings) + { + var config = new AmazonS3Config(); + + var serviceUrl = settings.S3ServiceUrl; + if (!string.IsNullOrEmpty(serviceUrl)) + { + config.ServiceURL = serviceUrl; + config.ForcePathStyle = true; // Required for S3-compatible endpoints (MinIO, LocalStack). + } + + var region = settings.S3Region; + if (!string.IsNullOrEmpty(region)) + { + if (!string.IsNullOrEmpty(serviceUrl)) + { + config.AuthenticationRegion = region; + } + else + { + config.RegionEndpoint = RegionEndpoint.GetBySystemName(region); + } + } + + // With no static keys the SDK's default credential chain resolves the ambient IAM role. + return !string.IsNullOrEmpty(settings.S3AccessKeyId) && !string.IsNullOrEmpty(settings.S3SecretAccessKey) + ? new AmazonS3Client(new BasicAWSCredentials(settings.S3AccessKeyId, settings.S3SecretAccessKey), config) + : new AmazonS3Client(config); + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/ExpectedLengthStream.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/ExpectedLengthStream.cs new file mode 100644 index 0000000000..914c66995c --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/ExpectedLengthStream.cs @@ -0,0 +1,65 @@ +namespace ServiceControl.Persistence.EFCore.Infrastructure; + +/// +/// Validates the byte count of a forward-only stream without buffering its contents. +/// +/// +/// Compressed body objects record their uncompressed length in metadata. Before cloud reads were +/// streamed, eager decompression verified that length before returning the body. This wrapper +/// preserves that integrity check at end of stream. A consumer that stops reading early does not +/// perform the final length check, which is necessary to keep response streaming and cancellation. +/// +public sealed class ExpectedLengthStream(Stream stream, long expectedLength) : Stream +{ + long bytesRead; + + public override bool CanRead => stream.CanRead; + public override bool CanSeek => false; + public override bool CanWrite => false; + public override long Length => throw new NotSupportedException(); + + public override long Position + { + get => bytesRead; + set => throw new NotSupportedException(); + } + + public override void Flush() => throw new NotSupportedException(); + + public override int Read(byte[] buffer, int offset, int count) => ValidateRead(stream.Read(buffer, offset, count)); + + public override int Read(Span buffer) => ValidateRead(stream.Read(buffer)); + + public override async Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken) => + ValidateRead(await stream.ReadAsync(buffer, offset, count, cancellationToken)); + + public override async ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) => + ValidateRead(await stream.ReadAsync(buffer, cancellationToken)); + + public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); + + public override void SetLength(long value) => throw new NotSupportedException(); + + public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException(); + + int ValidateRead(int count) + { + bytesRead += count; + if (bytesRead > expectedLength || (count == 0 && bytesRead != expectedLength)) + { + throw new InvalidOperationException("Decompressed body size does not match its metadata."); + } + + return count; + } + + protected override void Dispose(bool disposing) + { + if (disposing) + { + stream.Dispose(); + } + + base.Dispose(disposing); + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/MessageBodyFileResult.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/MessageBodyFileResult.cs index 531eba1f58..81ebf87a23 100644 --- a/src/ServiceControl.Persistence.EFCore/Infrastructure/MessageBodyFileResult.cs +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/MessageBodyFileResult.cs @@ -7,5 +7,4 @@ public class MessageBodyFileResult public required Stream Stream { get; set; } public required string ContentType { get; set; } public required int BodySize { get; set; } - public required string Etag { get; set; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/OwnedStream.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/OwnedStream.cs new file mode 100644 index 0000000000..be265364fb --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/OwnedStream.cs @@ -0,0 +1,68 @@ +namespace ServiceControl.Persistence.EFCore.Infrastructure; + +/// +/// Couples a returned stream to another resource that must remain alive while the stream is read. +/// +/// +/// S3 returns a response whose response stream is consumed later by ASP.NET Core, after the +/// persistence method has returned. Disposing this wrapper disposes both the exposed stream and +/// the owning response, releasing the underlying HTTP connection. +/// +public sealed class OwnedStream(Stream stream, IDisposable owner) : Stream +{ + public override bool CanRead => stream.CanRead; + public override bool CanSeek => stream.CanSeek; + public override bool CanWrite => stream.CanWrite; + public override long Length => stream.Length; + + public override long Position + { + get => stream.Position; + set => stream.Position = value; + } + + public override void Flush() => stream.Flush(); + + public override Task FlushAsync(CancellationToken cancellationToken) => stream.FlushAsync(cancellationToken); + + public override int Read(byte[] buffer, int offset, int count) => stream.Read(buffer, offset, count); + + public override int Read(Span buffer) => stream.Read(buffer); + + public override Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken) => + stream.ReadAsync(buffer, offset, count, cancellationToken); + + public override ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) => + stream.ReadAsync(buffer, cancellationToken); + + public override long Seek(long offset, SeekOrigin origin) => stream.Seek(offset, origin); + + public override void SetLength(long value) => stream.SetLength(value); + + public override void Write(byte[] buffer, int offset, int count) => stream.Write(buffer, offset, count); + + public override void Write(ReadOnlySpan buffer) => stream.Write(buffer); + + public override Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken) => + stream.WriteAsync(buffer, offset, count, cancellationToken); + + public override ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) => + stream.WriteAsync(buffer, cancellationToken); + + protected override void Dispose(bool disposing) + { + if (disposing) + { + try + { + stream.Dispose(); + } + finally + { + owner.Dispose(); + } + } + + base.Dispose(disposing); + } +} diff --git a/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj b/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj index f80428271a..5aca192107 100644 --- a/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj +++ b/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj @@ -19,6 +19,9 @@ + + + diff --git a/src/ServiceControl.Persistence.Tests.PostgreSql/ServiceControl.Persistence.Tests.PostgreSql.csproj b/src/ServiceControl.Persistence.Tests.PostgreSql/ServiceControl.Persistence.Tests.PostgreSql.csproj index 59de3364fa..3c63181a00 100644 --- a/src/ServiceControl.Persistence.Tests.PostgreSql/ServiceControl.Persistence.Tests.PostgreSql.csproj +++ b/src/ServiceControl.Persistence.Tests.PostgreSql/ServiceControl.Persistence.Tests.PostgreSql.csproj @@ -22,6 +22,8 @@ + + diff --git a/src/ServiceControl.Persistence.Tests.SqlServer/ServiceControl.Persistence.Tests.SqlServer.csproj b/src/ServiceControl.Persistence.Tests.SqlServer/ServiceControl.Persistence.Tests.SqlServer.csproj index 186518974d..c29c79b2e3 100644 --- a/src/ServiceControl.Persistence.Tests.SqlServer/ServiceControl.Persistence.Tests.SqlServer.csproj +++ b/src/ServiceControl.Persistence.Tests.SqlServer/ServiceControl.Persistence.Tests.SqlServer.csproj @@ -22,6 +22,8 @@ + + diff --git a/src/ServiceControl.Persistence.Tests/EFCore/AzureBlobBodyStorageTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/AzureBlobBodyStorageTests.cs new file mode 100644 index 0000000000..45c66a154a --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/AzureBlobBodyStorageTests.cs @@ -0,0 +1,142 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.IO; +using System.Text; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Implementation; +using Testcontainers.Azurite; + +[TestFixture] +[Platform(Exclude = "Win", Reason = "Azurite has no Windows container image")] +class AzureBlobBodyStorageTests +{ + AzuriteContainer azurite; + AzureBlobBodyStoragePersistence store; + + [OneTimeSetUp] + public async Task StartAzurite() + { + azurite = new AzuriteBuilder("mcr.microsoft.com/azure-storage/azurite:latest") + // We need this because the Azurite image is not yet compatible with the latest Azure SDK, + // see https://github.com/Azure/Azurite/issues/2623 + .WithCommand("--skipApiVersionCheck") + .Build(); + await azurite.StartAsync(); + } + + [OneTimeTearDown] + public async Task StopAzurite() + { + if (azurite != null) + { + await azurite.DisposeAsync(); + } + } + + [SetUp] + public async Task CreateContainer() + { + var settings = new TestSettings + { + ConnectionString = "not-used", + BodyStorageType = BodyStorageType.AzureBlob, + AzureBlobConnectionString = azurite.GetConnectionString(), + AzureBlobContainerName = $"bodies-{Guid.NewGuid():n}", + MinBodySizeForCompression = 64 + }; + + await new AzureBlobBodyStorageInstaller(settings).Provision(); + store = new AzureBlobBodyStoragePersistence(settings); + } + + [Test] + public async Task Round_trips_a_small_uncompressed_body() + { + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes("hello world"); + + await store.WriteBody(bodyId, body, "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.Multiple(() => + { + Assert.That(result.ContentType, Is.EqualTo("text/plain")); + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + }); + } + + [Test] + public async Task Round_trips_a_large_body_over_the_compression_threshold() + { + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes(new string('a', 100_000)); + + await store.WriteBody(bodyId, body, "application/json"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + } + + [Test] + public async Task Returns_null_for_a_missing_body() => + Assert.That(await store.ReadBody(Guid.NewGuid().ToString()), Is.Null); + + [Test] + public async Task Delete_removes_the_body() + { + var bodyId = Guid.NewGuid().ToString(); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("payload"), "text/plain"); + + await store.DeleteBody(bodyId); + + Assert.That(await store.ReadBody(bodyId), Is.Null); + } + + [Test] + public void Delete_of_a_missing_body_does_not_throw() => + Assert.DoesNotThrowAsync(() => store.DeleteBody(Guid.NewGuid().ToString())); + + [Test] + public async Task Rewriting_an_existing_body_keeps_the_first_write() + { + var bodyId = Guid.NewGuid().ToString(); + var original = Encoding.UTF8.GetBytes("original"); + + await store.WriteBody(bodyId, original, "text/plain"); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("different"), "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(original), "bodies are immutable, so the first write wins"); + } + } + + static byte[] ReadAll(Stream stream) + { + using var buffer = new MemoryStream(); + stream.CopyTo(buffer); + return buffer.ToArray(); + } + + sealed class TestSettings : EFPersisterSettings; +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs new file mode 100644 index 0000000000..bf85d372e7 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs @@ -0,0 +1,118 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.IO; +using System.Text; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Operations.BodyStorage; + +class BodyReadTests : ErrorIngestionTestBase +{ + const int Cap = 64; + + [SetUp] + public void ShrinkTheBodyCap() => EFSettings.MaxBodySizeToStore = Cap; + + [Test] + public async Task Fetches_an_inline_text_body() + { + var failure = new IngestedFailure { ContentType = "text/xml", Body = Encoding.UTF8.GetBytes("1") }; + await Ingest(failure); + + var result = await Fetch(failure.UniqueMessageIdString); + + Assert.That(result, Is.Not.Null); + using (Assert.EnterMultipleScope()) + { + Assert.That(result.HasResult, Is.True); + Assert.That(result.ContentType, Is.EqualTo("text/xml")); + Assert.That(Encoding.UTF8.GetString(ReadAll(result.Stream)), Is.EqualTo("1")); + } + } + + [Test] + public async Task Fetches_an_external_binary_body() + { + var body = BitConverter.GetBytes(0xDEADBEEF); + var failure = new IngestedFailure { ContentType = "application/octet-stream", Body = body }; + await Ingest(failure); + + var result = await Fetch(failure.UniqueMessageIdString); + + Assert.That(result, Is.Not.Null); + using (Assert.EnterMultipleScope()) + { + Assert.That(result.HasResult, Is.True); + Assert.That(result.ContentType, Is.EqualTo("application/octet-stream")); + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + } + + [Test] + public async Task Fetches_the_whole_body_for_large_text_not_the_inline_prefix() + { + var body = Encoding.UTF8.GetBytes(new string('a', Cap * 2)); + var failure = new IngestedFailure { Body = body }; + await Ingest(failure); + + var result = await Fetch(failure.UniqueMessageIdString); + + Assert.That(result, Is.Not.Null); + using (Assert.EnterMultipleScope()) + { + Assert.That(result.HasResult, Is.True); + Assert.That(ReadAll(result.Stream), Is.EqualTo(body), "external storage is authoritative, not the inline search prefix"); + } + } + + [Test] + public async Task Fetches_by_message_id() + { + var failure = new IngestedFailure { Body = Encoding.UTF8.GetBytes("1") }; + await Ingest(failure); + + var result = await Fetch(failure.MessageId); + + Assert.That(result, Is.Not.Null); + Assert.That(result.HasResult, Is.True); + } + + [Test] + public async Task Reports_no_body_for_an_empty_body() + { + var failure = new IngestedFailure { Body = [] }; + await Ingest(failure); + + var result = await Fetch(failure.UniqueMessageIdString); + + Assert.That(result, Is.Not.Null); + Assert.That(result.HasResult, Is.False); + } + + [Test] + public async Task Returns_null_for_an_unknown_message() + { + var result = await Fetch(Guid.NewGuid().ToString()); + + Assert.That(result, Is.Null); + } + + async Task Fetch(string bodyId) + { + using var scope = ServiceProvider.CreateScope(); + var bodyStorage = scope.ServiceProvider.GetRequiredService(); + return await bodyStorage.TryFetch(bodyId); + } + + static byte[] ReadAll(Stream stream) + { + using (stream) + { + using var buffer = new MemoryStream(); + stream.CopyTo(buffer); + return buffer.ToArray(); + } + } +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/BodyStoragePersistenceTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/BodyStoragePersistenceTests.cs new file mode 100644 index 0000000000..ff6c539450 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/BodyStoragePersistenceTests.cs @@ -0,0 +1,168 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.IO; +using System.Text; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Implementation; +using ServiceControl.Persistence.EFCore.Infrastructure; + +[TestFixture] +class BodyStoragePersistenceTests +{ + const string FileSystem = nameof(FileSystem); + const string InMemory = nameof(InMemory); + + string tempDir; + + [SetUp] + public void SetUp() => tempDir = Directory.CreateTempSubdirectory("sc-body-tests-").FullName; + + [TearDown] + public void TearDown() + { + try + { + Directory.Delete(tempDir, recursive: true); + } + catch (DirectoryNotFoundException) + { + } + } + + IBodyStoragePersistence CreateStore(string kind) => kind switch + { + InMemory => new InMemoryBodyStoragePersistence(), + FileSystem => new FileSystemBodyStoragePersistence(new TestSettings + { + ConnectionString = "not-used", + MessageBodyStoragePath = tempDir, + MinBodySizeForCompression = 64 + }), + _ => throw new ArgumentOutOfRangeException(nameof(kind)) + }; + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public async Task Round_trips_a_small_uncompressed_body(string kind) + { + var store = CreateStore(kind); + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes("hello world"); + + await store.WriteBody(bodyId, body, "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.Multiple(() => + { + Assert.That(result.ContentType, Is.EqualTo("text/plain")); + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + }); + } + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public async Task Round_trips_a_large_body_over_the_compression_threshold(string kind) + { + var store = CreateStore(kind); + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes(new string('a', 100_000)); + + await store.WriteBody(bodyId, body, "application/json"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + } + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public async Task Returns_null_for_a_missing_body(string kind) + { + var store = CreateStore(kind); + + Assert.That(await store.ReadBody(Guid.NewGuid().ToString()), Is.Null); + } + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public async Task Delete_removes_the_body(string kind) + { + var store = CreateStore(kind); + var bodyId = Guid.NewGuid().ToString(); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("payload"), "text/plain"); + + await store.DeleteBody(bodyId); + + Assert.That(await store.ReadBody(bodyId), Is.Null); + } + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public void Delete_of_a_missing_body_does_not_throw(string kind) + { + var store = CreateStore(kind); + + Assert.DoesNotThrowAsync(() => store.DeleteBody(Guid.NewGuid().ToString())); + } + + [TestCase(InMemory)] + [TestCase(FileSystem)] + public async Task Rewriting_an_existing_body_keeps_the_first_write(string kind) + { + var store = CreateStore(kind); + var bodyId = Guid.NewGuid().ToString(); + var original = Encoding.UTF8.GetBytes("original"); + + await store.WriteBody(bodyId, original, "text/plain"); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("different"), "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(original), "bodies are immutable, so the first write wins"); + } + } + + [Test] + public async Task Provisioning_creates_the_filesystem_storage_directory() + { + var path = Path.Combine(tempDir, "nested", "bodies"); + Assert.That(Directory.Exists(path), Is.False); + + await new FileSystemBodyStorageInstaller(new TestSettings + { + ConnectionString = "not-used", + BodyStorageType = BodyStorageType.FileSystem, + MessageBodyStoragePath = path + }).Provision(); + + Assert.That(Directory.Exists(path), Is.True); + } + + static byte[] ReadAll(Stream stream) + { + using var buffer = new MemoryStream(); + stream.CopyTo(buffer); + return buffer.ToArray(); + } + + sealed class TestSettings : EFPersisterSettings; +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/ErrorIngestionTestBase.cs b/src/ServiceControl.Persistence.Tests/EFCore/ErrorIngestionTestBase.cs index 32bc18e917..8d5e50e9f7 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/ErrorIngestionTestBase.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/ErrorIngestionTestBase.cs @@ -3,7 +3,6 @@ namespace ServiceControl.Persistence.Tests; using System; using System.Collections.Generic; using System.Linq; -using System.Threading; using System.Threading.Tasks; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; @@ -20,7 +19,7 @@ abstract class ErrorIngestionTestBase : PersistenceTestBase protected ErrorIngestionTestBase() => RegisterServices = services => services.AddSingleton(RecordedBodies); - protected RecordingBodyStoragePersistence RecordedBodies { get; } = new(); + protected InMemoryBodyStoragePersistence RecordedBodies { get; } = new(); protected EFPersisterSettings EFSettings => (EFPersisterSettings)PersistenceSettings; @@ -97,65 +96,4 @@ async Task Query(Func> query) return await query(dbContext); } - - protected class RecordingBodyStoragePersistence : IBodyStoragePersistence - { - readonly List written = []; - readonly List deleted = []; - - // Body ids whose deletion should throw, to exercise the sweep's tolerate-missing handling. - public HashSet FailDeleteFor { get; } = []; - - public IReadOnlyList Written - { - get - { - lock (written) - { - return [.. written]; - } - } - } - - public IReadOnlyList Deleted - { - get - { - lock (deleted) - { - return [.. deleted]; - } - } - } - - public Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) - { - lock (written) - { - written.Add(new StoredBody(bodyId, body.ToArray(), contentType)); - } - - return Task.CompletedTask; - } - - public Task ReadBody(string bodyId, CancellationToken cancellationToken = default) => - throw new NotImplementedException(); - - public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) - { - if (FailDeleteFor.Contains(bodyId)) - { - throw new InvalidOperationException($"Simulated missing body for {bodyId}"); - } - - lock (deleted) - { - deleted.Add(bodyId); - } - - return Task.CompletedTask; - } - - public record StoredBody(string BodyId, byte[] Body, string ContentType); - } } diff --git a/src/ServiceControl.Persistence.Tests/EFCore/ExpectedLengthStreamTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/ExpectedLengthStreamTests.cs new file mode 100644 index 0000000000..67570b2add --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/ExpectedLengthStreamTests.cs @@ -0,0 +1,40 @@ +namespace ServiceControl.Persistence.Tests; + +using System.IO; +using NUnit.Framework; +using ServiceControl.Persistence.EFCore.Infrastructure; + +[TestFixture] +class ExpectedLengthStreamTests +{ + [Test] + public void Allows_a_stream_with_the_expected_length() + { + using var stream = new ExpectedLengthStream(new MemoryStream([1, 2, 3]), 3); + + Assert.That(ReadAll(stream), Is.EqualTo(new byte[] { 1, 2, 3 })); + } + + [Test] + public void Throws_when_a_stream_ends_before_the_expected_length() + { + using var stream = new ExpectedLengthStream(new MemoryStream([1, 2]), 3); + + Assert.That(() => ReadAll(stream), Throws.InvalidOperationException); + } + + [Test] + public void Throws_when_a_stream_exceeds_the_expected_length() + { + using var stream = new ExpectedLengthStream(new MemoryStream([1, 2, 3]), 2); + + Assert.That(() => ReadAll(stream), Throws.InvalidOperationException); + } + + static byte[] ReadAll(Stream stream) + { + using var buffer = new MemoryStream(); + stream.CopyTo(buffer); + return buffer.ToArray(); + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs b/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs new file mode 100644 index 0000000000..f3b2c7c0cc --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs @@ -0,0 +1,91 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Collections.Generic; +using System.IO; +using System.Threading; +using System.Threading.Tasks; +using ServiceControl.Persistence.EFCore.Infrastructure; + +class InMemoryBodyStoragePersistence : IBodyStoragePersistence +{ + readonly object gate = new(); + readonly List written = []; + readonly List deleted = []; + readonly Dictionary store = []; + + public HashSet FailDeleteFor { get; } = []; + + public IReadOnlyList Written + { + get + { + lock (gate) + { + return [.. written]; + } + } + } + + public IReadOnlyList Deleted + { + get + { + lock (gate) + { + return [.. deleted]; + } + } + } + + public Task WriteBody(string bodyId, ReadOnlyMemory body, string contentType, CancellationToken cancellationToken = default) + { + var entry = new StoredBody(bodyId, body.ToArray(), contentType); + + lock (gate) + { + written.Add(entry); + store.TryAdd(bodyId, entry); // First write wins, like the real stores. + } + + return Task.CompletedTask; + } + + public Task ReadBody(string bodyId, CancellationToken cancellationToken = default) + { + StoredBody entry; + lock (gate) + { + entry = store.GetValueOrDefault(bodyId); + } + + var result = entry is null + ? null + : new MessageBodyFileResult + { + Stream = new MemoryStream(entry.Body, writable: false), + ContentType = entry.ContentType, + BodySize = entry.Body.Length + }; + + return Task.FromResult(result); + } + + public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) + { + if (FailDeleteFor.Contains(bodyId)) + { + throw new InvalidOperationException($"Simulated missing body for {bodyId}"); + } + + lock (gate) + { + deleted.Add(bodyId); + store.Remove(bodyId); + } + + return Task.CompletedTask; + } + + public record StoredBody(string BodyId, byte[] Body, string ContentType); +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/S3BodyStorageTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/S3BodyStorageTests.cs new file mode 100644 index 0000000000..53de694113 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/S3BodyStorageTests.cs @@ -0,0 +1,142 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.IO; +using System.Text; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.Persistence.EFCore.Abstractions; +using ServiceControl.Persistence.EFCore.Implementation; +using Testcontainers.LocalStack; + +[TestFixture] +[Platform(Exclude = "Win", Reason = "LocalStack has no Windows container image")] +class S3BodyStorageTests +{ + LocalStackContainer localStack; + S3BodyStoragePersistence store; + + [OneTimeSetUp] + public async Task StartLocalStack() + { + // We need to pin the tag to 4.4.0 because after that version LocalStack has become a paid service for commercial use. + localStack = new LocalStackBuilder("localstack/localstack:4.4.0").Build(); + await localStack.StartAsync(); + } + + [OneTimeTearDown] + public async Task StopLocalStack() + { + if (localStack != null) + { + await localStack.DisposeAsync(); + } + } + + [SetUp] + public async Task CreateBucket() + { + var settings = new TestSettings + { + ConnectionString = "not-used", + BodyStorageType = BodyStorageType.S3, + S3ServiceUrl = localStack.GetConnectionString(), + S3Region = "us-east-1", + S3AccessKeyId = "test", + S3SecretAccessKey = "test", + S3BucketName = $"bodies-{Guid.NewGuid():n}", + MinBodySizeForCompression = 64 + }; + + await new S3BodyStorageInstaller(settings).Provision(); + store = new S3BodyStoragePersistence(settings); + } + + [Test] + public async Task Round_trips_a_small_uncompressed_body() + { + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes("hello world"); + + await store.WriteBody(bodyId, body, "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.Multiple(() => + { + Assert.That(result.ContentType, Is.EqualTo("text/plain")); + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + }); + } + + [Test] + public async Task Round_trips_a_large_body_over_the_compression_threshold() + { + var bodyId = Guid.NewGuid().ToString(); + var body = Encoding.UTF8.GetBytes(new string('a', 100_000)); + + await store.WriteBody(bodyId, body, "application/json"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(body)); + } + + Assert.That(result.BodySize, Is.EqualTo(body.Length)); + } + + [Test] + public async Task Returns_null_for_a_missing_body() => + Assert.That(await store.ReadBody(Guid.NewGuid().ToString()), Is.Null); + + [Test] + public async Task Delete_removes_the_body() + { + var bodyId = Guid.NewGuid().ToString(); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("payload"), "text/plain"); + + await store.DeleteBody(bodyId); + + Assert.That(await store.ReadBody(bodyId), Is.Null); + } + + [Test] + public void Delete_of_a_missing_body_does_not_throw() => + Assert.DoesNotThrowAsync(() => store.DeleteBody(Guid.NewGuid().ToString())); + + [Test] + public async Task Rewriting_an_existing_body_keeps_the_first_write() + { + var bodyId = Guid.NewGuid().ToString(); + var original = Encoding.UTF8.GetBytes("original"); + + await store.WriteBody(bodyId, original, "text/plain"); + await store.WriteBody(bodyId, Encoding.UTF8.GetBytes("different"), "text/plain"); + + var result = await store.ReadBody(bodyId); + + Assert.That(result, Is.Not.Null); + using (result.Stream) + { + Assert.That(ReadAll(result.Stream), Is.EqualTo(original), "bodies are immutable, so the first write wins"); + } + } + + static byte[] ReadAll(Stream stream) + { + using var buffer = new MemoryStream(); + stream.CopyTo(buffer); + return buffer.ToArray(); + } + + sealed class TestSettings : EFPersisterSettings; +} diff --git a/src/ServiceControl.Persistence/IBodyStorageInstaller.cs b/src/ServiceControl.Persistence/IBodyStorageInstaller.cs new file mode 100644 index 0000000000..c16abbd6be --- /dev/null +++ b/src/ServiceControl.Persistence/IBodyStorageInstaller.cs @@ -0,0 +1,9 @@ +namespace ServiceControl.Persistence; + +using System.Threading; +using System.Threading.Tasks; + +public interface IBodyStorageInstaller +{ + Task Provision(CancellationToken cancellationToken = default); +} diff --git a/src/ServiceControl/Hosting/Commands/SetupCommand.cs b/src/ServiceControl/Hosting/Commands/SetupCommand.cs index 8cfd79d05e..fd3d8ebbe2 100644 --- a/src/ServiceControl/Hosting/Commands/SetupCommand.cs +++ b/src/ServiceControl/Hosting/Commands/SetupCommand.cs @@ -55,6 +55,11 @@ public override async Task Execute(HostArguments args, Settings settings) { await databaseMigrator.ApplyMigrations(); } + + if (scope.ServiceProvider.GetService() is { } bodyStorageInstaller) + { + await bodyStorageInstaller.Provision(); + } } await host.StopAsync();