Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
113 changes: 113 additions & 0 deletions src/SmtpServer.Benchmarks/DataStoreBenchmarks.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
using System;
using System.Buffers;
using System.IO;
using System.IO.Pipelines;
using System.Threading;
using System.Threading.Tasks;
using BenchmarkDotNet.Attributes;
using MimeKit;
using SmtpServer.ComponentModel;
using SmtpServer.Protocol;
using SmtpServer.Storage;
using SmtpClient = MailKit.Net.Smtp.SmtpClient;

namespace SmtpServer.Benchmarks
{
[MemoryDiagnoser]
[ShortRunJob]
public class DataStoreBenchmarks
{
const int Port = 9026;

SmtpServer _smtpServer;
CancellationTokenSource _smtpServerCancellationTokenSource;
SmtpClient _smtpClient;
MimeMessage _message;

public enum StoreMode
{
BufferedMaterializing,
StreamingDrain
}

[Params(StoreMode.BufferedMaterializing, StoreMode.StreamingDrain)]
public StoreMode Mode { get; set; }

[GlobalSetup]
public void SmtpServerSetup()
{
_message = MimeMessage.Load(typeof(DataStoreBenchmarks).Assembly.GetManifestResourceStream("SmtpServer.Benchmarks.Test3.eml"));
_smtpServerCancellationTokenSource = new CancellationTokenSource();

var serviceProvider = new ServiceProvider();
serviceProvider.Add(Mode == StoreMode.StreamingDrain
? (IMessageStore)new StreamingDrainMessageStore()
: new BufferedMaterializingMessageStore());

_smtpServer = new SmtpServer(
new SmtpServerOptionsBuilder()
.Port(Port, false)
.Build(),
serviceProvider);

_ = _smtpServer.StartAsync(_smtpServerCancellationTokenSource.Token);

_smtpClient = new SmtpClient();
_smtpClient.Connect("localhost", Port);
}

[GlobalCleanup]
public Task SmtpServerCleanupAsync()
{
_smtpClient.Disconnect(true);
_smtpClient.Dispose();

_smtpServerCancellationTokenSource.Cancel();
_smtpServerCancellationTokenSource.Dispose();

return _smtpServer.ShutdownTask;
}

[Benchmark]
public void SendMessage()
{
_smtpClient.Send(_message);
}

sealed class BufferedMaterializingMessageStore : MessageStore
{
public override Task<SmtpResponse> SaveAsync(ISessionContext context, IMessageTransaction transaction, ReadOnlySequence<byte> buffer, CancellationToken cancellationToken)
{
_ = buffer.ToArray();

return Task.FromResult(SmtpResponse.Ok);
}
}

sealed class StreamingDrainMessageStore : MessageStore, IStreamingMessageStore
{
public override Task<SmtpResponse> SaveAsync(ISessionContext context, IMessageTransaction transaction, ReadOnlySequence<byte> buffer, CancellationToken cancellationToken)
{
throw new NotSupportedException();
}

public async Task<SmtpResponse> SaveAsync(ISessionContext context, IMessageTransaction transaction, PipeReader reader, CancellationToken cancellationToken)
{
while (true)
{
var result = await reader.ReadAsync(cancellationToken).ConfigureAwait(false);
var buffer = result.Buffer;

reader.AdvanceTo(buffer.End);

if (result.IsCompleted)
{
break;
}
}

return SmtpResponse.Ok;
}
}
}
}
18 changes: 7 additions & 11 deletions src/SmtpServer.Benchmarks/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,17 +8,13 @@ public class Program
{
public static void Main(string[] args)
{
//var summary = BenchmarkRunner.Run<TokenizerBenchmarks>(
// ManualConfig
// .Create(DefaultConfig.Instance)
// .With(ConfigOptions.DisableOptimizationsValidator));

//var summary = BenchmarkRunner.Run<ThroughputBenchmarks>();

var summary = BenchmarkRunner.Run<ThroughputBenchmarks>(
ManualConfig
.Create(DefaultConfig.Instance)
.With(ConfigOptions.DisableOptimizationsValidator));
BenchmarkSwitcher
.FromAssembly(typeof(Program).Assembly)
.Run(
args,
ManualConfig
.Create(DefaultConfig.Instance)
.WithOptions(ConfigOptions.DisableOptimizationsValidator));
}
}
}
62 changes: 62 additions & 0 deletions src/SmtpServer.Tests/PipeReaderTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using System.Text;
using System.Threading.Tasks;
using SmtpServer.IO;
using SmtpServer.Protocol;
using SmtpServer.Text;
using Xunit;

Expand Down Expand Up @@ -93,5 +94,66 @@ await reader.ReadDotBlockAsync(
// assert
Assert.Equal("abcd\r\n.1234", text);
}

[Fact]
public async Task CanStreamBlockWithDotStuffingRemoved()
{
// arrange
var reader = CreatePipeReader("abcd\r\n..1234\r\n.\r\n");
var writer = new Pipe();

var maxMessageSizeOptions = new MaxMessageSizeOptions();

// act
await reader.ReadDotBlockAsync(writer.Writer, maxMessageSizeOptions);
var text = await ReadAllAsync(writer.Reader);

// assert
Assert.Equal("abcd\r\n.1234", text);
}

[Fact]
public async Task CanEnforceMaxMessageSizeWhenStreamingBlock()
{
// arrange
var reader = CreatePipeReader("abcd\r\n1234\r\n.\r\n");
var writer = new Pipe();

var maxMessageSizeOptions = new MaxMessageSizeOptions(MaxMessageSizeHandling.Strict, 5);

// act
var exception = await Assert.ThrowsAsync<SmtpResponseException>(
async () => await reader.ReadDotBlockAsync(writer.Writer, maxMessageSizeOptions));

// assert
Assert.True(exception.IsQuitRequested);
}

static async Task<string> ReadAllAsync(PipeReader reader)
{
using var stream = new MemoryStream();

while (true)
{
var result = await reader.ReadAsync();
var buffer = result.Buffer;

foreach (var segment in buffer)
{
stream.Write(segment.Span);
}

reader.AdvanceTo(buffer.End);

if (result.IsCompleted)
{
break;
}
}

reader.Complete();

return Encoding.ASCII.GetString(stream.ToArray());
}
}
}
62 changes: 62 additions & 0 deletions src/SmtpServer.Tests/SmtpServerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,10 @@
using SmtpServer.Storage;
using SmtpServer.Tests.Mocks;
using System;
using System.Buffers;
using System.Diagnostics;
using System.IO;
using System.IO.Pipelines;
using System.Linq;
using System.Net;
using System.Net.Security;
Expand Down Expand Up @@ -52,6 +54,21 @@ public void CanReceiveMessage()
}
}

[Fact]
public void CanReceiveMessageUsingStreamingMessageStore()
{
var streamingMessageStore = new StreamingMockMessageStore();

using (CreateServer(services => services.Add(streamingMessageStore)))
{
MailClient.Send(MailClient.Message(from: "test1@test.com", to: "test2@test.com", text: "streamed body"));
}

Assert.True(streamingMessageStore.StreamingSaveCalled);
Assert.False(streamingMessageStore.BufferedSaveCalled);
Assert.Contains("streamed body", streamingMessageStore.Message);
}

[Theory]
[InlineData("Assunto teste acento çãõáéíóú", "utf-8")]
[InlineData("שלום שלום שלום", "windows-1255")]
Expand Down Expand Up @@ -655,5 +672,50 @@ SmtpServerDisposable CreateServer(
/// The cancellation token source for the test.
/// </summary>
public CancellationTokenSource CancellationTokenSource { get; }

sealed class StreamingMockMessageStore : MessageStore, IStreamingMessageStore
{
public override Task<SmtpResponse> SaveAsync(ISessionContext context, IMessageTransaction transaction, ReadOnlySequence<byte> buffer, CancellationToken cancellationToken)
{
BufferedSaveCalled = true;

return Task.FromResult(SmtpResponse.Ok);
}

public async Task<SmtpResponse> SaveAsync(ISessionContext context, IMessageTransaction transaction, PipeReader reader, CancellationToken cancellationToken)
{
StreamingSaveCalled = true;

using var stream = new MemoryStream();

while (true)
{
var result = await reader.ReadAsync(cancellationToken);
var buffer = result.Buffer;

foreach (var segment in buffer)
{
stream.Write(segment.Span);
}

reader.AdvanceTo(buffer.End);

if (result.IsCompleted)
{
break;
}
}

Message = Encoding.UTF8.GetString(stream.ToArray());

return SmtpResponse.Ok;
}

public bool BufferedSaveCalled { get; private set; }

public bool StreamingSaveCalled { get; private set; }

public string Message { get; private set; }
}
}
}
Loading