Skip to content
Merged
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
1 change: 1 addition & 0 deletions Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
<PackageVersion Include="Microsoft.Extensions.Diagnostics" Version="8.0.1" />
<PackageVersion Include="Microsoft.Extensions.Hosting.Abstractions" Version="8.0.1" />
<PackageVersion Include="Microsoft.Extensions.Options.ConfigurationExtensions" Version="8.0.0" />
<PackageVersion Include="Microsoft.Extensions.Diagnostics.Testing" Version="8.10.0" />
<PackageVersion Include="Microsoft.Extensions.TimeProvider.Testing" Version="8.10.0" />
<PackageVersion Include="Microsoft.NET.Test.Sdk" Version="17.14.1" />
<PackageVersion Include="Microsoft.SourceLink.GitHub" Version="8.0.0" />
Expand Down
327 changes: 0 additions & 327 deletions OutboxKit.sln

This file was deleted.

40 changes: 40 additions & 0 deletions OutboxKit.slnx
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
<Solution>
<Configurations>
<Platform Name="Any CPU"/>
<Platform Name="x64"/>
<Platform Name="x86"/>
</Configurations>
<Folder Name="/build/"/>
<Folder Name="/samples/"/>
<Folder Name="/samples/mongodb/">
<Project Path="samples/mongodb/MongoDbMultiDbPollingSample/MongoDbMultiDbPollingSample.csproj"/>
<Project Path="samples/mongodb/MongoDbPollingSample/MongoDbPollingSample.csproj"/>
</Folder>
<Folder Name="/samples/mysql/">
<Project Path="samples/mysql/MySqlEfMultiDbPollingSample/MySqlEfMultiDbPollingSample.csproj"/>
<Project Path="samples/mysql/MySqlEfPollingSample/MySqlEfPollingSample.csproj"/>
</Folder>
<Folder Name="/samples/mysql/MySqlEndToEndPollingSample/">
<Project Path="samples/mysql/MySqlEndToEndPollingSample/Consumer/Consumer.csproj"/>
<Project Path="samples/mysql/MySqlEndToEndPollingSample/OutOfProcessProducer/OutOfProcessProducer.csproj"/>
<Project Path="samples/mysql/MySqlEndToEndPollingSample/Producer/Producer.csproj"/>
<Project Path="samples/mysql/MySqlEndToEndPollingSample/ProducerShared/ProducerShared.csproj"/>
</Folder>
<Folder Name="/samples/postgresql/">
<Project Path="samples/postgresql/PostgreSqlEfMultiDbPollingSample/PostgreSqlEfMultiDbPollingSample.csproj"/>
<Project Path="samples/postgresql/PostgreSqlEfPollingSample/PostgreSqlEfPollingSample.csproj"/>
</Folder>
<Folder Name="/src/">
<Project Path="src/Core.OpenTelemetry/Core.OpenTelemetry.csproj"/>
<Project Path="src/Core/Core.csproj"/>
<Project Path="src/MongoDb/MongoDb.csproj"/>
<Project Path="src/MySql/MySql.csproj"/>
<Project Path="src/PostgreSql/PostgreSql.csproj"/>
</Folder>
<Folder Name="/tests/">
<Project Path="tests/Core.Tests/Core.Tests.csproj"/>
<Project Path="tests/MongoDb.Tests/MongoDb.Tests.csproj"/>
<Project Path="tests/MySql.Tests/MySql.Tests.csproj"/>
<Project Path="tests/PostgreSql.Tests/PostgreSql.Tests.csproj"/>
</Folder>
</Solution>
2 changes: 1 addition & 1 deletion build/CakeRunner.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
#:sdk Cake.Sdk

const string solutionPath = "./OutboxKit.sln";
const string solutionPath = "./OutboxKit.slnx";
const string librariesPath = "./src/";
const string artifactsPath = "./artifacts/";

Expand Down
4 changes: 2 additions & 2 deletions docker-compose.yml
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
services:
mongodb:
image: mongo
image: mongo:8
Comment thread
joaofbantunes marked this conversation as resolved.
container_name: mongodb
command: ["--replSet", "rs0", "--bind_ip_all"]
ports:
- "27017:27017"

init-mongo:
image: mongo
image: mongo:8
container_name: init-mongo
depends_on:
- mongodb
Expand Down
2 changes: 1 addition & 1 deletion global.json
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,6 @@
"rollForward": "latestFeature"
},
"msbuild-sdks": {
"Cake.Sdk": "6.1.1"
"Cake.Sdk": "6.2.0"
}
}
5 changes: 4 additions & 1 deletion src/MongoDb/Polling/ConfigurationImplementation.cs
Original file line number Diff line number Diff line change
Expand Up @@ -161,12 +161,15 @@ public void ConfigureServices(OutboxKey key, IServiceCollection services)

services.AddKeyedSingleton(key, _dbFactory);

services.AddKeyedSingleton<ChangeStreamListener>(key);

services.AddKeyedSingleton<DistributedLockThingy>(
key,
(s, _) => ActivatorUtilities.CreateInstance<DistributedLockThingy>(
s,
_lockBaseSettings,
s.GetRequiredKeyedService<Func<OutboxKey, IServiceProvider, IMongoDatabase>>(key)(key, s)));
s.GetRequiredKeyedService<Func<OutboxKey, IServiceProvider, IMongoDatabase>>(key)(key, s),
s.GetRequiredKeyedService<ChangeStreamListener>(key)));

_collectionConfigurator.ConfigureMe(new GetMongoDbOutboxCollectionConfigured(
key,
Expand Down
144 changes: 122 additions & 22 deletions src/MongoDb/Synchronization/ChangeStreamListener.cs
Original file line number Diff line number Diff line change
@@ -1,16 +1,17 @@
using Microsoft.Extensions.Logging;
using MongoDB.Driver;
using Nito.AsyncEx;

namespace YakShaveFx.OutboxKit.MongoDb.Synchronization;

internal sealed class ChangeStreamListener(
AsyncAutoResetEvent autoResetEvent,
IChangeStreamCursor<ChangeStreamDocument<DistributedLockDocument>> cursor,
CancellationTokenSource cts) : IAsyncDisposable
internal interface IChangeStreamNotifier : IAsyncDisposable
{
public Task WaitAsync() => autoResetEvent.WaitAsync(cts.Token);
Task OnChangeAsync(CancellationToken ct);
}

public static async Task<ChangeStreamListener> StartAsync(
internal sealed partial class ChangeStreamListener(ILogger<ChangeStreamListener> logger)
{
public async Task<IChangeStreamNotifier> ListenAsync(
IMongoCollection<DistributedLockDocument> collection,
DistributedLockDefinition lockDefinition,
CancellationToken ct)
Expand All @@ -21,41 +22,140 @@ public static async Task<ChangeStreamListener> StartAsync(
.Match(d => d.DocumentKey["_id"] == lockDefinition.Id),
new ChangeStreamOptions
{
BatchSize = 1,
MaxAwaitTime = TimeSpan.FromMinutes(5)
BatchSize = 1
},
ct);

var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
var autoResetEvent = new AsyncAutoResetEvent();
_ = Task.Run(async () =>
var backgroundTask = Task.Run(async () =>
{
while (!cts.Token.IsCancellationRequested && await cursor.MoveNextAsync(cts.Token))
try
{
// only yield when a change is detected, don't care about the amount, just that there is a change
if (cursor.Current.Any())
while (!cts.Token.IsCancellationRequested && await cursor.MoveNextAsync(cts.Token))
{
autoResetEvent.Set();
// only yield when a relevant change is detected, don't care about the amount, just that there is a relevant change
if (cursor.Current.Any(d => ShouldYield(d, lockDefinition, logger)))
{
autoResetEvent.Set();
}
}
}
catch (OperationCanceledException) when (cts.IsCancellationRequested)
{
// expected when the cancellation token is canceled
}
catch (Exception ex)
{
// log and forget is only acceptable here,
// because there always is a parallel process relying on delays to double-check things
LogErrorWatchingForLockChanges(logger, ex, lockDefinition.Id, lockDefinition.Context);
}
}, cts.Token);

var watcher = new ChangeStreamListener(autoResetEvent, cursor, cts);
var watcher = new ChangeStreamNotifier(autoResetEvent, cursor, cts, backgroundTask);
return watcher;

/*
* we care if:
* - any delete to the lock document while it should be up
* - any insert or replace to the lock document while it should be up, but only if the owner is different
* (otherwise we'd get notified by what the current lock is doing)
*/
static bool ShouldYield(
ChangeStreamDocument<DistributedLockDocument> document,
DistributedLockDefinition lockDefinition,
ILogger logger)
{
if (document.OperationType is ChangeStreamOperationType.Delete)
{
LogLockDeletion(logger, lockDefinition.Id, lockDefinition.Context);
return true;
}

if (document.OperationType is ChangeStreamOperationType.Insert or ChangeStreamOperationType.Replace
&& document.FullDocument.Owner != lockDefinition.Owner)
{
LogLockChangeWithDifferentOwner(
logger,
document.OperationType,
lockDefinition.Owner,
document.FullDocument.Owner,
lockDefinition.Id,
lockDefinition.Context);

return true;
}

LogIrrelevantLockChange(logger, document.OperationType, lockDefinition.Id, lockDefinition.Context);

return false;
}
}

public async ValueTask DisposeAsync()
private sealed class ChangeStreamNotifier(
AsyncAutoResetEvent autoResetEvent,
IChangeStreamCursor<ChangeStreamDocument<DistributedLockDocument>> cursor,
CancellationTokenSource cts,
Task backgroundTask) : IChangeStreamNotifier
{
try
public async Task OnChangeAsync(CancellationToken ct)
{
await cts.CancelAsync();
using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(ct, cts.Token);
await autoResetEvent.WaitAsync(linkedCts.Token);
}
catch (Exception)

public async ValueTask DisposeAsync()
{
// try to cancel, but don't throw if it fails
}
try
{
await cts.CancelAsync();
await backgroundTask;
}
catch (Exception)
{
// try to cancel and await the task, but don't throw if it fails
}

cts.Dispose();
cursor.Dispose();
cursor.Dispose();
cts.Dispose();
}
}

[LoggerMessage(
LogLevel.Debug,
Message = "Lock deletion detected (id \"{Id}\" context \"{Context}\")")]
private static partial void LogLockDeletion(ILogger logger, string? id, string? context);


[LoggerMessage(
LogLevel.Debug,
Message =
"Lock change detected with different owner (operation \"{OperationType}\" expected owner \"{ExpectedOwner}\" actual owner \"{ActualOwner}\" id \"{Id}\" context \"{Context}\")")]
private static partial void LogLockChangeWithDifferentOwner(
ILogger logger,
ChangeStreamOperationType operationType,
string expectedOwner,
string? actualOwner,
string id,
string? context);

[LoggerMessage(
LogLevel.Debug,
Message = "Irrelevant lock change detected (operation \"{OperationType}\" id \"{Id}\" context \"{Context}\")")]
private static partial void LogIrrelevantLockChange(
ILogger logger,
ChangeStreamOperationType operationType,
string id,
string? context);

[LoggerMessage(
LogLevel.Warning,
Message =
"An error occurred while watching for lock changes, falling back to time based alternatives (id \"{Id}\" context \"{Context}\")")]
private static partial void LogErrorWatchingForLockChanges(
ILogger logger,
Exception ex,
string id,
string? context);
}
30 changes: 20 additions & 10 deletions src/MongoDb/Synchronization/DistributedLockThingy.cs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ internal sealed partial class DistributedLockThingy(
DistributedLockSettings settings,
IMongoDatabase database,
TimeProvider timeProvider,
ChangeStreamListener changeStreamListener,
ILogger<DistributedLockThingy> logger)
{
private readonly IMongoCollection<DistributedLockDocument> _collection =
Expand All @@ -28,7 +29,7 @@ public async Task<IDistributedLock> AcquireAsync(DistributedLockDefinition lockD
}
catch (Exception)
{
await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
throw;
}
}
Expand All @@ -42,7 +43,7 @@ public async Task<IDistributedLock> AcquireAsync(DistributedLockDefinition lockD
{
if (!await InnerTryAcquireAsync(internalLockDefinition, ct))
{
await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
return null;
}

Expand All @@ -52,7 +53,7 @@ public async Task<IDistributedLock> AcquireAsync(DistributedLockDefinition lockD
}
catch (Exception)
{
await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
throw;
}
}
Expand All @@ -61,14 +62,14 @@ private async Task<InternalDistributedLockDefinition> CreateInternalLockDefiniti
DistributedLockDefinition lockDefinition,
CancellationToken ct)
{
var changeStreamListener = _changeStreamsEnabled
? await ChangeStreamListener.StartAsync(_collection, lockDefinition, ct)
var changeStreamNotifier = _changeStreamsEnabled
? await changeStreamListener.ListenAsync(_collection, lockDefinition, ct)
: null;

return new InternalDistributedLockDefinition
{
Definition = lockDefinition,
ChangeStreamListener = changeStreamListener
ChangeStreamNotifier = changeStreamNotifier
};
}

Expand Down Expand Up @@ -182,7 +183,7 @@ private async Task WatchAndKeepTryingToAcquireAsync(InternalDistributedLockDefin
while (!ct.IsCancellationRequested)
{
// listener is not null when change streams are enabled
await lockDefinition.ChangeStreamListener!.WaitAsync();
await lockDefinition.ChangeStreamNotifier!.OnChangeAsync(ct);
if (await InnerTryAcquireAsync(lockDefinition, ct)) return;
}
}
Expand Down Expand Up @@ -227,9 +228,11 @@ private void KickoffKeepAlive(InternalDistributedLockDefinition lockDefinition,
try
{
var delayTask = Task.Delay(keepAliveInterval, timeProvider, linkedTokenSource.Token);

// listener is not null when change streams are enabled
await Task.WhenAny(delayTask, lockDefinition.ChangeStreamListener!.WaitAsync());
watchLockLossTask = lockDefinition.ChangeStreamNotifier!.OnChangeAsync(linkedTokenSource.Token);

await Task.WhenAny(delayTask, watchLockLossTask);

if (!delayTask.IsCompleted)
{
Expand Down Expand Up @@ -333,6 +336,13 @@ private sealed class DistributedLock(
CancellationTokenSource keepAliveCts,
Func<InternalDistributedLockDefinition, CancellationTokenSource, ValueTask> releaseLock) : IDistributedLock
{
public ValueTask DisposeAsync() => releaseLock(definition, keepAliveCts);
public async ValueTask DisposeAsync()
{
await releaseLock(definition, keepAliveCts);
if (definition.ChangeStreamNotifier is not null)
{
await definition.ChangeStreamNotifier.DisposeAsync();
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Synchronization;
internal sealed class InternalDistributedLockDefinition
{
public required DistributedLockDefinition Definition { get; init; }
public required ChangeStreamListener? ChangeStreamListener { get; init; }
public required IChangeStreamNotifier? ChangeStreamNotifier { get; init; }

public string Id => Definition.Id;
public string Owner => Definition.Owner;
Expand Down
1 change: 1 addition & 0 deletions tests/MongoDb.Tests/MongoDb.Tests.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
<PackageReference Include="FluentAssertions"/>
<PackageReference Include="Microsoft.Extensions.Diagnostics.Testing"/>
<PackageReference Include="Microsoft.Extensions.TimeProvider.Testing"/>
<PackageReference Include="Microsoft.NET.Test.Sdk"/>
<PackageReference Include="xunit.v3"/>
Expand Down
Loading
Loading