Skip to content

Commit 6f33fbf

Browse files
authored
Merge pull request #5646 from Particular/john/body_storage
Implement configurable message body storage
2 parents 784e566 + 3a7d242 commit 6f33fbf

49 files changed

Lines changed: 2137 additions & 112 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

src/Directory.Packages.props

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,9 +6,11 @@
66
<ItemGroup Label="Versions for direct package references">
77
<PackageVersion Include="Autofac" Version="9.3.0" />
88
<PackageVersion Include="AWSSDK.CloudWatch" Version="4.0.10.5" />
9+
<PackageVersion Include="AWSSDK.S3" Version="4.0.101.4" />
910
<PackageVersion Include="Azure.Identity" Version="1.21.0" />
1011
<PackageVersion Include="Azure.ResourceManager.Monitor" Version="1.3.1" />
1112
<PackageVersion Include="Azure.ResourceManager.ServiceBus" Version="1.1.0" />
13+
<PackageVersion Include="Azure.Storage.Blobs" Version="12.27.0" />
1214
<PackageVersion Include="ByteSize" Version="2.1.2" />
1315
<PackageVersion Include="Caliburn.Micro" Version="5.0.258" />
1416
<PackageVersion Include="DnsClient" Version="1.8.0" />
@@ -76,6 +78,8 @@
7678
<PackageVersion Include="PropertyChanging.Fody" Version="1.31.0" />
7779
<PackageVersion Include="PublicApiGenerator" Version="11.5.4" />
7880
<PackageVersion Include="RavenDB.Embedded" Version="6.2.17" />
81+
<PackageVersion Include="Testcontainers.Azurite" Version="4.13.0" />
82+
<PackageVersion Include="Testcontainers.LocalStack" Version="4.13.0" />
7983
<PackageVersion Include="Testcontainers.MsSql" Version="4.13.0" />
8084
<PackageVersion Include="Testcontainers.PostgreSql" Version="4.13.0" />
8185
<PackageVersion Include="ReactiveUI.WPF" Version="22.3.1" />

src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistence.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ public void AddPersistence(IServiceCollection services)
1212
{
1313
RegisterSettings(services);
1414
ConfigureDbContext(services);
15-
RegisterDataStores(services);
15+
RegisterDataStores(services, settings);
1616

1717
services.AddSingleton<IIngestionSqlDialect, PostgreSqlIngestionSqlDialect>();
1818
}
@@ -23,6 +23,7 @@ public void AddInstaller(IServiceCollection services)
2323
ConfigureDbContext(services);
2424

2525
services.AddScoped<IDatabaseMigrator, PostgreSqlDatabaseMigrator>();
26+
RegisterBodyStorageInstaller(services, settings);
2627
}
2728

2829
void RegisterSettings(IServiceCollection services)

src/ServiceControl.Persistence.EFCore.PostgreSql/PostgreSqlPersistenceConfiguration.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,6 @@ public override IPersistence Create(PersistenceSettings settings)
1515
return new PostgreSqlPersistence((PostgreSqlPersisterSettings)settings);
1616
}
1717

18-
protected override EFPersisterSettings CreateSettings(string connectionString) =>
19-
new PostgreSqlPersisterSettings { ConnectionString = connectionString };
18+
protected override EFPersisterSettings CreateSettings(string connectionString, BodyStorageSettings bodyStorage) =>
19+
new PostgreSqlPersisterSettings { ConnectionString = connectionString, BodyStorage = bodyStorage };
2020
}

src/ServiceControl.Persistence.EFCore.PostgreSql/persistence.manifest

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,58 @@
1313
"Name": "ServiceControl/Database/CommandTimeout",
1414
"Mandatory": false
1515
},
16+
{
17+
"Name": "ServiceControl/MessageBody/StorageType",
18+
"Mandatory": true
19+
},
1620
{
1721
"Name": "ServiceControl/MessageBody/StoragePath",
1822
"Mandatory": false
1923
},
24+
{
25+
"Name": "ServiceControl/MessageBody/Azure/ConnectionString",
26+
"Mandatory": false
27+
},
28+
{
29+
"Name": "ServiceControl/MessageBody/Azure/ServiceUri",
30+
"Mandatory": false
31+
},
32+
{
33+
"Name": "ServiceControl/MessageBody/Azure/ManagedIdentityClientId",
34+
"Mandatory": false
35+
},
36+
{
37+
"Name": "ServiceControl/MessageBody/Azure/AuthorityHost",
38+
"Mandatory": false
39+
},
40+
{
41+
"Name": "ServiceControl/MessageBody/Azure/ContainerName",
42+
"Mandatory": false
43+
},
44+
{
45+
"Name": "ServiceControl/MessageBody/S3/BucketName",
46+
"Mandatory": false
47+
},
48+
{
49+
"Name": "ServiceControl/MessageBody/S3/KeyPrefix",
50+
"Mandatory": false
51+
},
52+
{
53+
"Name": "ServiceControl/MessageBody/S3/Region",
54+
"Mandatory": false
55+
},
56+
{
57+
"Name": "ServiceControl/MessageBody/S3/ServiceUrl",
58+
"Mandatory": false
59+
},
60+
{
61+
"Name": "ServiceControl/MessageBody/S3/AccessKeyId",
62+
"Mandatory": false
63+
},
64+
{
65+
"Name": "ServiceControl/MessageBody/S3/SecretAccessKey",
66+
"Mandatory": false
67+
},
2068
{
2169
"Name": "ServiceControl/MessageBody/MinCompressionSize",
2270
"Mandatory": false

src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistence.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ public void AddPersistence(IServiceCollection services)
1212
{
1313
RegisterSettings(services);
1414
ConfigureDbContext(services);
15-
RegisterDataStores(services);
15+
RegisterDataStores(services, settings);
1616

1717
services.AddSingleton<IIngestionSqlDialect, SqlServerIngestionSqlDialect>();
1818
}
@@ -23,6 +23,7 @@ public void AddInstaller(IServiceCollection services)
2323
ConfigureDbContext(services);
2424

2525
services.AddScoped<IDatabaseMigrator, SqlServerDatabaseMigrator>();
26+
RegisterBodyStorageInstaller(services, settings);
2627
}
2728

2829
void RegisterSettings(IServiceCollection services)

src/ServiceControl.Persistence.EFCore.SqlServer/SqlServerPersistenceConfiguration.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,6 @@ public override IPersistence Create(PersistenceSettings settings)
1515
return new SqlServerPersistence((SqlServerPersisterSettings)settings);
1616
}
1717

18-
protected override EFPersisterSettings CreateSettings(string connectionString) =>
19-
new SqlServerPersisterSettings { ConnectionString = connectionString };
18+
protected override EFPersisterSettings CreateSettings(string connectionString, BodyStorageSettings bodyStorage) =>
19+
new SqlServerPersisterSettings { ConnectionString = connectionString, BodyStorage = bodyStorage };
2020
}

src/ServiceControl.Persistence.EFCore.SqlServer/persistence.manifest

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,58 @@
1313
"Name": "ServiceControl/Database/CommandTimeout",
1414
"Mandatory": false
1515
},
16+
{
17+
"Name": "ServiceControl/MessageBody/StorageType",
18+
"Mandatory": true
19+
},
1620
{
1721
"Name": "ServiceControl/MessageBody/StoragePath",
1822
"Mandatory": false
1923
},
24+
{
25+
"Name": "ServiceControl/MessageBody/Azure/ConnectionString",
26+
"Mandatory": false
27+
},
28+
{
29+
"Name": "ServiceControl/MessageBody/Azure/ServiceUri",
30+
"Mandatory": false
31+
},
32+
{
33+
"Name": "ServiceControl/MessageBody/Azure/ManagedIdentityClientId",
34+
"Mandatory": false
35+
},
36+
{
37+
"Name": "ServiceControl/MessageBody/Azure/AuthorityHost",
38+
"Mandatory": false
39+
},
40+
{
41+
"Name": "ServiceControl/MessageBody/Azure/ContainerName",
42+
"Mandatory": false
43+
},
44+
{
45+
"Name": "ServiceControl/MessageBody/S3/BucketName",
46+
"Mandatory": false
47+
},
48+
{
49+
"Name": "ServiceControl/MessageBody/S3/KeyPrefix",
50+
"Mandatory": false
51+
},
52+
{
53+
"Name": "ServiceControl/MessageBody/S3/Region",
54+
"Mandatory": false
55+
},
56+
{
57+
"Name": "ServiceControl/MessageBody/S3/ServiceUrl",
58+
"Mandatory": false
59+
},
60+
{
61+
"Name": "ServiceControl/MessageBody/S3/AccessKeyId",
62+
"Mandatory": false
63+
},
64+
{
65+
"Name": "ServiceControl/MessageBody/S3/SecretAccessKey",
66+
"Mandatory": false
67+
},
2068
{
2169
"Name": "ServiceControl/MessageBody/MinCompressionSize",
2270
"Mandatory": false
Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,28 @@
1+
namespace ServiceControl.Persistence.EFCore.Abstractions;
2+
3+
public sealed class AzureBlobBodyStorageSettings : BodyStorageSettings
4+
{
5+
public const string DefaultContainerName = "error-bodies";
6+
7+
public required AzureBlobAuthentication Authentication { get; set; }
8+
public string ContainerName { get; set; } = DefaultContainerName;
9+
}
10+
11+
// Shared-key and managed-identity auth are mutually exclusive, and the managed identity options are
12+
// meaningless alongside a connection string.
13+
public abstract class AzureBlobAuthentication;
14+
15+
public sealed class AzureBlobSharedKeyAuthentication : AzureBlobAuthentication
16+
{
17+
public required string ConnectionString { get; set; }
18+
}
19+
20+
public sealed class AzureBlobManagedIdentityAuthentication : AzureBlobAuthentication
21+
{
22+
public required Uri ServiceUri { get; set; }
23+
public string? ClientId { get; set; }
24+
25+
// Steers the login endpoint for sovereign clouds; when unset the SDK honours the
26+
// AZURE_AUTHORITY_HOST environment variable.
27+
public Uri? AuthorityHost { get; set; }
28+
}

src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs

Lines changed: 48 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,12 @@
11
namespace ServiceControl.Persistence.EFCore.Abstractions;
22

33
using Microsoft.Extensions.DependencyInjection;
4+
using Microsoft.Extensions.DependencyInjection.Extensions;
45
using NServiceBus.Unicast.Subscriptions.MessageDrivenSubscriptions;
56
using Particular.LicensingComponent.Persistence;
67
using ServiceControl.Operations.BodyStorage;
78
using ServiceControl.Persistence.EFCore.Implementation;
9+
using ServiceControl.Persistence.EFCore.Implementation.BodyStorage;
810
using ServiceControl.Persistence.EFCore.Implementation.UnitOfWork;
911
using ServiceControl.Persistence.EFCore.Infrastructure;
1012
using ServiceControl.Persistence.MessageRedirects;
@@ -13,7 +15,7 @@ namespace ServiceControl.Persistence.EFCore.Abstractions;
1315

1416
public abstract class BasePersistence
1517
{
16-
protected static void RegisterDataStores(IServiceCollection services)
18+
protected static void RegisterDataStores(IServiceCollection services, EFPersisterSettings settings)
1719
{
1820
services.AddSingleton(TimeProvider.System);
1921
services.AddSingleton<MinimumRequiredStorageState>();
@@ -51,6 +53,50 @@ protected static void RegisterDataStores(IServiceCollection services)
5153

5254
services.AddSingleton<ILicensingDataStore, LicensingDataStore>();
5355

54-
services.AddSingleton<IBodyStoragePersistence, FakeBodyStoragePersistence>();
56+
RegisterBodyStorage(services, settings);
57+
}
58+
59+
// Settings are registered under their concrete type so each store resolves only what it can act on.
60+
static void RegisterBodyStorage(IServiceCollection services, EFPersisterSettings settings)
61+
{
62+
switch (settings.BodyStorage)
63+
{
64+
case FileSystemBodyStorageSettings fileSystem:
65+
services.TryAddSingleton(fileSystem);
66+
services.AddSingleton<IBodyStoragePersistence, FileSystemBodyStoragePersistence>();
67+
break;
68+
case AzureBlobBodyStorageSettings azureBlob:
69+
services.TryAddSingleton(azureBlob);
70+
services.AddSingleton<IBodyStoragePersistence, AzureBlobBodyStoragePersistence>();
71+
break;
72+
case S3BodyStorageSettings s3:
73+
services.TryAddSingleton(s3);
74+
services.AddSingleton<IBodyStoragePersistence, S3BodyStoragePersistence>();
75+
break;
76+
default:
77+
throw new ArgumentOutOfRangeException(nameof(settings), settings.BodyStorage, "Unknown body storage type.");
78+
}
79+
}
80+
81+
// Only stores needing setup-time provisioning register an installer; SetupCommand skips when none is.
82+
protected static void RegisterBodyStorageInstaller(IServiceCollection services, EFPersisterSettings settings)
83+
{
84+
switch (settings.BodyStorage)
85+
{
86+
case FileSystemBodyStorageSettings fileSystem:
87+
services.TryAddSingleton(fileSystem);
88+
services.AddScoped<IBodyStorageInstaller, FileSystemBodyStorageInstaller>();
89+
break;
90+
case AzureBlobBodyStorageSettings azureBlob:
91+
services.TryAddSingleton(azureBlob);
92+
services.AddScoped<IBodyStorageInstaller, AzureBlobBodyStorageInstaller>();
93+
break;
94+
case S3BodyStorageSettings s3:
95+
services.TryAddSingleton(s3);
96+
services.AddScoped<IBodyStorageInstaller, S3BodyStorageInstaller>();
97+
break;
98+
default:
99+
throw new ArgumentOutOfRangeException(nameof(settings), settings.BodyStorage, "Unknown body storage type.");
100+
}
55101
}
56102
}
Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
namespace ServiceControl.Persistence.EFCore.Abstractions;
2+
3+
/// <summary>
4+
/// Settings for the selected body storage type.
5+
/// </summary>
6+
/// <remarks>
7+
/// One subclass per storage type, so a store only ever receives the settings it can act on and the
8+
/// configuration layer's validation is carried by the type rather than re-asserted at the point of use.
9+
/// </remarks>
10+
public abstract class BodyStorageSettings
11+
{
12+
public const int DefaultMinCompressionSize = 4096;
13+
public const int DefaultMaxBodySizeToStore = 102400; // 100 kb
14+
15+
public int MinCompressionSize { get; set; } = DefaultMinCompressionSize;
16+
public int MaxBodySizeToStore { get; set; } = DefaultMaxBodySizeToStore;
17+
}

0 commit comments

Comments
 (0)