diff --git a/MailerQ.Test/MessageStorageTest.cs b/MailerQ.Test/MessageStorageTest.cs index d72f82ba..2c62b3cb 100644 --- a/MailerQ.Test/MessageStorageTest.cs +++ b/MailerQ.Test/MessageStorageTest.cs @@ -71,6 +71,7 @@ public void Constructor_should_throw_exception_with_not_implemented_supported_en [InlineData(SimpleValidMongoUri)] [InlineData("mongodb://mongos1.example.com,mongos2.example.com/?readPreference=secondary")] [InlineData("mongodb://mongos1.example.com,mongos2.example.com/database?readPreference=secondary")] + [InlineData("s3://accesskey:secretkey@region/bucketname")] [Theory] public void Constructor_should_create_instance_with_valid_supported_engine_uri(string uri) { diff --git a/MailerQ/MailerQ.csproj b/MailerQ/MailerQ.csproj index 063d9507..e5ac31ce 100644 --- a/MailerQ/MailerQ.csproj +++ b/MailerQ/MailerQ.csproj @@ -15,6 +15,7 @@ + diff --git a/MailerQ/MessageStore/IMessageStorage.cs b/MailerQ/MessageStore/IMessageStorage.cs index 66fba475..a865c3ed 100644 --- a/MailerQ/MessageStore/IMessageStorage.cs +++ b/MailerQ/MessageStore/IMessageStorage.cs @@ -17,6 +17,11 @@ public interface IMessageStorage /// Task InsertAsync(string message, int secondsToExpire = DefaultSecondsToExpire, CancellationToken cancellationToken = default); + /// + /// The default storage name to use if not indicated explicitly + /// + const string DefaultStorageName = "mailerq"; + /// /// Message time to live into the storage engine before expire. Default 25 hours /// diff --git a/MailerQ/MessageStore/MessageStorage.cs b/MailerQ/MessageStore/MessageStorage.cs index bdcaa49c..bdae3be7 100644 --- a/MailerQ/MessageStore/MessageStorage.cs +++ b/MailerQ/MessageStore/MessageStorage.cs @@ -38,13 +38,12 @@ public MessageStorage(IOptions options) private IMessageStorage CreateConcretStorage(StorageEngines storageEngine, string uri) { - switch (storageEngine) + return storageEngine switch { - case StorageEngines.MongoDB: - return new MongoDBMessageStorage(uri); - default: - throw new NotImplementedException($"{storageEngine} message storage engine is not implement yet."); - } + StorageEngines.MongoDB => new MongoDBMessageStorage(uri), + StorageEngines.S3 => new S3MessageStorage(uri), + _ => throw new NotImplementedException($"{storageEngine} message storage engine is not implement yet."), + }; } /// diff --git a/MailerQ/MessageStore/MongoDBMessageStorage.cs b/MailerQ/MessageStore/MongoDBMessageStorage.cs index 4484f211..32e0e82e 100644 --- a/MailerQ/MessageStore/MongoDBMessageStorage.cs +++ b/MailerQ/MessageStore/MongoDBMessageStorage.cs @@ -10,7 +10,6 @@ namespace MailerQ.MessageStore internal class MongoDBMessageStorage : IMessageStorage { private const int MessageMaxSuppportedSize = 15728640; // 15 MB - private const string DataBase = "mailerq"; private readonly string Collection = "mime"; private readonly IMongoCollection messages; @@ -18,7 +17,7 @@ public MongoDBMessageStorage(string url) { var mongoUrl = MongoUrl.Create(url); var mongoClient = new MongoClient(mongoUrl); - var database = mongoClient.GetDatabase(mongoUrl.DatabaseName ?? DataBase); + var database = mongoClient.GetDatabase(mongoUrl.DatabaseName ?? IMessageStorage.DefaultStorageName); messages = database.GetCollection(Collection); } diff --git a/MailerQ/MessageStore/S3MessageStorage.cs b/MailerQ/MessageStore/S3MessageStorage.cs new file mode 100644 index 00000000..2245711b --- /dev/null +++ b/MailerQ/MessageStore/S3MessageStorage.cs @@ -0,0 +1,35 @@ +using Amazon.S3; +using Amazon.S3.Model; +using System; +using System.Threading; +using System.Threading.Tasks; + +namespace MailerQ.MessageStore +{ + internal class S3MessageStorage : IMessageStorage + { + private readonly IAmazonS3 _client; + + public S3MessageStorage(string url) + { + var config = new AmazonS3Config { ServiceURL = url }; + _client = new AmazonS3Client(config); + } + + public async Task InsertAsync(string message, int secondsToExpire, CancellationToken cancellationToken) + { + var key = Guid.NewGuid().ToString(); + var request = new PutObjectRequest + { + BucketName = IMessageStorage.DefaultStorageName, + Key = key, + ContentType = "text/plain", + ContentBody = message, + }; + + await _client.PutObjectAsync(request, cancellationToken); + + return key; + } + } +}