Skip to content

Commit

Permalink
Added xml comments to OutboxSweeper
Browse files Browse the repository at this point in the history
  • Loading branch information
preardon committed Feb 20, 2022
1 parent 43045f2 commit bd3fbb3
Show file tree
Hide file tree
Showing 2 changed files with 14 additions and 20 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ private void DoWork(object state)
IAmACommandProcessor commandProcessor = scope.ServiceProvider.GetService<IAmACommandProcessor>();

var outBoxSweeper = new OutboxSweeper(
milliSecondsSinceSent: _options.MinimumMessageAge,
millisecondsSinceSent: _options.MinimumMessageAge,
commandProcessor: commandProcessor,
_options.BatchSize,
_options.UseBulk);
Expand Down
32 changes: 13 additions & 19 deletions src/Paramore.Brighter/OutboxSweeper.cs
Original file line number Diff line number Diff line change
@@ -1,48 +1,42 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

namespace Paramore.Brighter
namespace Paramore.Brighter
{
public class OutboxSweeper
{
private readonly int _milliSecondsSinceSent;
private readonly int _millisecondsSinceSent;
private readonly IAmACommandProcessor _commandProcessor;
private readonly int _batchSize;
private readonly bool _useBulk;

/// <summary>
/// This sweeper clears an outbox of any outstanding messages within the time interval
/// </summary>
/// <param name="milliSecondsSinceSent">How long can a message sit in the box before we attempt to resend</param>
/// <param name="millisecondsSinceSent">How long can a message sit in the box before we attempt to resend</param>
/// <param name="commandProcessor">Who should post the messages</param>
/// <param name="batchSize">The maximum number of messages to dispatch.</param>
/// <param name="useBulk">Use the producers bulk dispatch functionality.</param>
public OutboxSweeper(int milliSecondsSinceSent, IAmACommandProcessor commandProcessor, int batchSize = 100,
public OutboxSweeper(int millisecondsSinceSent, IAmACommandProcessor commandProcessor, int batchSize = 100,
bool useBulk = false)
{
_milliSecondsSinceSent = milliSecondsSinceSent;
_millisecondsSinceSent = millisecondsSinceSent;
_commandProcessor = commandProcessor;
_batchSize = batchSize;
_useBulk = useBulk;
}

/// <summary>
/// Dispatches the oldest un-dispatched messages from the outbox in a background thread.
/// </summary>
public void Sweep()
{
_commandProcessor.ClearOutbox(_batchSize, _milliSecondsSinceSent);
}

public Task SweepAsync(CancellationToken cancellationToken = default)
{
_commandProcessor.ClearAsyncOutbox(_batchSize, _milliSecondsSinceSent, _useBulk);

return Task.CompletedTask;
_commandProcessor.ClearOutbox(_batchSize, _millisecondsSinceSent);
}

/// <summary>
/// Dispatches the oldest un-dispatched messages from the asynchronous outbox in a background thread.
/// </summary>
public void SweepAsyncOutbox()
{
_commandProcessor.ClearAsyncOutbox(_batchSize, _milliSecondsSinceSent, _useBulk);
_commandProcessor.ClearAsyncOutbox(_batchSize, _millisecondsSinceSent, _useBulk);
}
}
}

0 comments on commit bd3fbb3

Please sign in to comment.