diff --git a/Directory.Packages.props b/Directory.Packages.props
index eb861c5..3ac3bcc 100644
--- a/Directory.Packages.props
+++ b/Directory.Packages.props
@@ -15,6 +15,7 @@
+
diff --git a/OutboxKit.sln b/OutboxKit.sln
deleted file mode 100644
index 2ef2179..0000000
--- a/OutboxKit.sln
+++ /dev/null
@@ -1,327 +0,0 @@
-
-Microsoft Visual Studio Solution File, Format Version 12.00
-# Visual Studio Version 17
-VisualStudioVersion = 17.0.31903.59
-MinimumVisualStudioVersion = 10.0.40219.1
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "src", "src", "{E1748C96-94CD-479F-BF71-585EB56C9FBF}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Core", "src\Core\Core.csproj", "{9E17490D-8407-441F-AF33-620EAD39C9AD}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "samples", "samples", "{2C240B41-C6C9-4921-9139-ECB2AA468316}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MySql", "src\MySql\MySql.csproj", "{16BB39E2-046B-4390-8375-37BCC15F81EF}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "tests", "tests", "{7DDB9CD3-11D6-4036-B0F5-19315D4FA1DD}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MySql.Tests", "tests\MySql.Tests\MySql.Tests.csproj", "{F566873C-2663-4B99-A6B9-259FCD6F8A2B}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Core.Tests", "tests\Core.Tests\Core.Tests.csproj", "{95240203-AFB6-4404-BBB9-8EB1F30E91F8}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "build", "build", "{A2D93BAA-65E6-4AF6-A2A4-9D53F74B5977}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Core.OpenTelemetry", "src\Core.OpenTelemetry\Core.OpenTelemetry.csproj", "{C4A3B901-3625-46B8-A224-DD4DEE1A5E65}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "mysql", "mysql", "{65461D72-9FB5-444E-B668-91F7C7C448A3}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MySqlEfMultiDbPollingSample", "samples\mysql\MySqlEfMultiDbPollingSample\MySqlEfMultiDbPollingSample.csproj", "{B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MySqlEfPollingSample", "samples\mysql\MySqlEfPollingSample\MySqlEfPollingSample.csproj", "{570E3464-9ABD-44E5-A1E8-D896BF779540}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "MySqlEndToEndPollingSample", "MySqlEndToEndPollingSample", "{846CDA36-9F33-4870-8C3B-0BC8DC3EC11B}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Producer", "samples\mysql\MySqlEndToEndPollingSample\Producer\Producer.csproj", "{D799F4D6-7170-493F-85DD-5999FDFF1085}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Consumer", "samples\mysql\MySqlEndToEndPollingSample\Consumer\Consumer.csproj", "{68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "OutOfProcessProducer", "samples\mysql\MySqlEndToEndPollingSample\OutOfProcessProducer\OutOfProcessProducer.csproj", "{EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ProducerShared", "samples\mysql\MySqlEndToEndPollingSample\ProducerShared\ProducerShared.csproj", "{13B80164-E50C-448E-9A38-C6D7A5313DC9}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MongoDb", "src\MongoDb\MongoDb.csproj", "{62F7D497-5170-4B0D-80B6-38C597FB979C}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "mongodb", "mongodb", "{114820B8-6912-4F74-85F6-EB25CAE48BF9}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MongoDbPollingSample", "samples\mongodb\MongoDbPollingSample\MongoDbPollingSample.csproj", "{C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MongoDbMultiDbPollingSample", "samples\mongodb\MongoDbMultiDbPollingSample\MongoDbMultiDbPollingSample.csproj", "{F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "MongoDb.Tests", "tests\MongoDb.Tests\MongoDb.Tests.csproj", "{1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PostgreSql", "src\PostgreSql\PostgreSql.csproj", "{D21719E8-67F7-4A2D-94E9-3B528471732C}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PostgreSql.Tests", "tests\PostgreSql.Tests\PostgreSql.Tests.csproj", "{C33EBB60-B465-4114-A764-4BDD78739136}"
-EndProject
-Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "postgresql", "postgresql", "{832FBD7C-83CB-926A-3032-625B67293BF7}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PostgreSqlEfPollingSample", "samples\postgresql\PostgreSqlEfPollingSample\PostgreSqlEfPollingSample.csproj", "{30EFE166-34B4-48FA-AE61-F2E463A8800B}"
-EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "PostgreSqlEfMultiDbPollingSample", "samples\postgresql\PostgreSqlEfMultiDbPollingSample\PostgreSqlEfMultiDbPollingSample.csproj", "{756BCA88-515D-4B59-BA11-3A9357D084CA}"
-EndProject
-Global
- GlobalSection(SolutionConfigurationPlatforms) = preSolution
- Debug|Any CPU = Debug|Any CPU
- Debug|x64 = Debug|x64
- Debug|x86 = Debug|x86
- Release|Any CPU = Release|Any CPU
- Release|x64 = Release|x64
- Release|x86 = Release|x86
- EndGlobalSection
- GlobalSection(ProjectConfigurationPlatforms) = postSolution
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|x64.ActiveCfg = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|x64.Build.0 = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|x86.ActiveCfg = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Debug|x86.Build.0 = Debug|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|Any CPU.Build.0 = Release|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|x64.ActiveCfg = Release|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|x64.Build.0 = Release|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|x86.ActiveCfg = Release|Any CPU
- {9E17490D-8407-441F-AF33-620EAD39C9AD}.Release|x86.Build.0 = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|x64.ActiveCfg = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|x64.Build.0 = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|x86.ActiveCfg = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Debug|x86.Build.0 = Debug|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|Any CPU.Build.0 = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|x64.ActiveCfg = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|x64.Build.0 = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|x86.ActiveCfg = Release|Any CPU
- {16BB39E2-046B-4390-8375-37BCC15F81EF}.Release|x86.Build.0 = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|x64.ActiveCfg = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|x64.Build.0 = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|x86.ActiveCfg = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Debug|x86.Build.0 = Debug|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|Any CPU.Build.0 = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|x64.ActiveCfg = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|x64.Build.0 = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|x86.ActiveCfg = Release|Any CPU
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B}.Release|x86.Build.0 = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|x64.ActiveCfg = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|x64.Build.0 = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|x86.ActiveCfg = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Debug|x86.Build.0 = Debug|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|Any CPU.Build.0 = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|x64.ActiveCfg = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|x64.Build.0 = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|x86.ActiveCfg = Release|Any CPU
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8}.Release|x86.Build.0 = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|x64.ActiveCfg = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|x64.Build.0 = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|x86.ActiveCfg = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Debug|x86.Build.0 = Debug|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|Any CPU.Build.0 = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|x64.ActiveCfg = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|x64.Build.0 = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|x86.ActiveCfg = Release|Any CPU
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65}.Release|x86.Build.0 = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|x64.ActiveCfg = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|x64.Build.0 = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|x86.ActiveCfg = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Debug|x86.Build.0 = Debug|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|Any CPU.Build.0 = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|x64.ActiveCfg = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|x64.Build.0 = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|x86.ActiveCfg = Release|Any CPU
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60}.Release|x86.Build.0 = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|x64.ActiveCfg = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|x64.Build.0 = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|x86.ActiveCfg = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Debug|x86.Build.0 = Debug|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|Any CPU.Build.0 = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|x64.ActiveCfg = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|x64.Build.0 = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|x86.ActiveCfg = Release|Any CPU
- {570E3464-9ABD-44E5-A1E8-D896BF779540}.Release|x86.Build.0 = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|x64.ActiveCfg = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|x64.Build.0 = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|x86.ActiveCfg = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Debug|x86.Build.0 = Debug|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|Any CPU.Build.0 = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|x64.ActiveCfg = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|x64.Build.0 = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|x86.ActiveCfg = Release|Any CPU
- {D799F4D6-7170-493F-85DD-5999FDFF1085}.Release|x86.Build.0 = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|x64.ActiveCfg = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|x64.Build.0 = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|x86.ActiveCfg = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Debug|x86.Build.0 = Debug|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|Any CPU.Build.0 = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|x64.ActiveCfg = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|x64.Build.0 = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|x86.ActiveCfg = Release|Any CPU
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0}.Release|x86.Build.0 = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|x64.ActiveCfg = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|x64.Build.0 = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|x86.ActiveCfg = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Debug|x86.Build.0 = Debug|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|Any CPU.Build.0 = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|x64.ActiveCfg = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|x64.Build.0 = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|x86.ActiveCfg = Release|Any CPU
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC}.Release|x86.Build.0 = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|x64.ActiveCfg = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|x64.Build.0 = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|x86.ActiveCfg = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Debug|x86.Build.0 = Debug|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|Any CPU.Build.0 = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|x64.ActiveCfg = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|x64.Build.0 = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|x86.ActiveCfg = Release|Any CPU
- {13B80164-E50C-448E-9A38-C6D7A5313DC9}.Release|x86.Build.0 = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|x64.ActiveCfg = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|x64.Build.0 = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|x86.ActiveCfg = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Debug|x86.Build.0 = Debug|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|Any CPU.Build.0 = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|x64.ActiveCfg = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|x64.Build.0 = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|x86.ActiveCfg = Release|Any CPU
- {62F7D497-5170-4B0D-80B6-38C597FB979C}.Release|x86.Build.0 = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|x64.ActiveCfg = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|x64.Build.0 = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|x86.ActiveCfg = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Debug|x86.Build.0 = Debug|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|Any CPU.Build.0 = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|x64.ActiveCfg = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|x64.Build.0 = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|x86.ActiveCfg = Release|Any CPU
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874}.Release|x86.Build.0 = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|x64.ActiveCfg = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|x64.Build.0 = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|x86.ActiveCfg = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Debug|x86.Build.0 = Debug|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|Any CPU.Build.0 = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|x64.ActiveCfg = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|x64.Build.0 = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|x86.ActiveCfg = Release|Any CPU
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D}.Release|x86.Build.0 = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|x64.ActiveCfg = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|x64.Build.0 = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|x86.ActiveCfg = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Debug|x86.Build.0 = Debug|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|Any CPU.Build.0 = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|x64.ActiveCfg = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|x64.Build.0 = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|x86.ActiveCfg = Release|Any CPU
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C}.Release|x86.Build.0 = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|x64.ActiveCfg = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|x64.Build.0 = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|x86.ActiveCfg = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Debug|x86.Build.0 = Debug|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|Any CPU.Build.0 = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|x64.ActiveCfg = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|x64.Build.0 = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|x86.ActiveCfg = Release|Any CPU
- {D21719E8-67F7-4A2D-94E9-3B528471732C}.Release|x86.Build.0 = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|x64.ActiveCfg = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|x64.Build.0 = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|x86.ActiveCfg = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Debug|x86.Build.0 = Debug|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|Any CPU.Build.0 = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|x64.ActiveCfg = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|x64.Build.0 = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|x86.ActiveCfg = Release|Any CPU
- {C33EBB60-B465-4114-A764-4BDD78739136}.Release|x86.Build.0 = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|x64.ActiveCfg = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|x64.Build.0 = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|x86.ActiveCfg = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Debug|x86.Build.0 = Debug|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|Any CPU.Build.0 = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|x64.ActiveCfg = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|x64.Build.0 = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|x86.ActiveCfg = Release|Any CPU
- {30EFE166-34B4-48FA-AE61-F2E463A8800B}.Release|x86.Build.0 = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|x64.ActiveCfg = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|x64.Build.0 = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|x86.ActiveCfg = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Debug|x86.Build.0 = Debug|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|Any CPU.Build.0 = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|x64.ActiveCfg = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|x64.Build.0 = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|x86.ActiveCfg = Release|Any CPU
- {756BCA88-515D-4B59-BA11-3A9357D084CA}.Release|x86.Build.0 = Release|Any CPU
- EndGlobalSection
- GlobalSection(SolutionProperties) = preSolution
- HideSolutionNode = FALSE
- EndGlobalSection
- GlobalSection(NestedProjects) = preSolution
- {9E17490D-8407-441F-AF33-620EAD39C9AD} = {E1748C96-94CD-479F-BF71-585EB56C9FBF}
- {16BB39E2-046B-4390-8375-37BCC15F81EF} = {E1748C96-94CD-479F-BF71-585EB56C9FBF}
- {F566873C-2663-4B99-A6B9-259FCD6F8A2B} = {7DDB9CD3-11D6-4036-B0F5-19315D4FA1DD}
- {95240203-AFB6-4404-BBB9-8EB1F30E91F8} = {7DDB9CD3-11D6-4036-B0F5-19315D4FA1DD}
- {C4A3B901-3625-46B8-A224-DD4DEE1A5E65} = {E1748C96-94CD-479F-BF71-585EB56C9FBF}
- {65461D72-9FB5-444E-B668-91F7C7C448A3} = {2C240B41-C6C9-4921-9139-ECB2AA468316}
- {B6AA2DEE-8319-4A6F-A8CC-4E5B6F145B60} = {65461D72-9FB5-444E-B668-91F7C7C448A3}
- {570E3464-9ABD-44E5-A1E8-D896BF779540} = {65461D72-9FB5-444E-B668-91F7C7C448A3}
- {846CDA36-9F33-4870-8C3B-0BC8DC3EC11B} = {65461D72-9FB5-444E-B668-91F7C7C448A3}
- {D799F4D6-7170-493F-85DD-5999FDFF1085} = {846CDA36-9F33-4870-8C3B-0BC8DC3EC11B}
- {68B9A32E-A1BF-4E8B-9C7A-8EE7BF968FC0} = {846CDA36-9F33-4870-8C3B-0BC8DC3EC11B}
- {EFEE7AA2-19CC-4029-BCEE-2C79E148C8EC} = {846CDA36-9F33-4870-8C3B-0BC8DC3EC11B}
- {13B80164-E50C-448E-9A38-C6D7A5313DC9} = {846CDA36-9F33-4870-8C3B-0BC8DC3EC11B}
- {62F7D497-5170-4B0D-80B6-38C597FB979C} = {E1748C96-94CD-479F-BF71-585EB56C9FBF}
- {114820B8-6912-4F74-85F6-EB25CAE48BF9} = {2C240B41-C6C9-4921-9139-ECB2AA468316}
- {C097DEE3-EDF3-4CD4-A7F2-DD4796A41874} = {114820B8-6912-4F74-85F6-EB25CAE48BF9}
- {F5C31DA0-FE7B-4471-884E-E8BEE095ED6D} = {114820B8-6912-4F74-85F6-EB25CAE48BF9}
- {1AE7C2D3-AC2D-4FD9-98B7-1256350A704C} = {7DDB9CD3-11D6-4036-B0F5-19315D4FA1DD}
- {D21719E8-67F7-4A2D-94E9-3B528471732C} = {E1748C96-94CD-479F-BF71-585EB56C9FBF}
- {C33EBB60-B465-4114-A764-4BDD78739136} = {7DDB9CD3-11D6-4036-B0F5-19315D4FA1DD}
- {832FBD7C-83CB-926A-3032-625B67293BF7} = {2C240B41-C6C9-4921-9139-ECB2AA468316}
- {30EFE166-34B4-48FA-AE61-F2E463A8800B} = {832FBD7C-83CB-926A-3032-625B67293BF7}
- {756BCA88-515D-4B59-BA11-3A9357D084CA} = {832FBD7C-83CB-926A-3032-625B67293BF7}
- EndGlobalSection
-EndGlobal
diff --git a/OutboxKit.slnx b/OutboxKit.slnx
new file mode 100644
index 0000000..cfd7bb4
--- /dev/null
+++ b/OutboxKit.slnx
@@ -0,0 +1,40 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/build/CakeRunner.cs b/build/CakeRunner.cs
index d5bf9de..01a635f 100644
--- a/build/CakeRunner.cs
+++ b/build/CakeRunner.cs
@@ -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/";
diff --git a/docker-compose.yml b/docker-compose.yml
index f9cb557..5f893de 100644
--- a/docker-compose.yml
+++ b/docker-compose.yml
@@ -1,13 +1,13 @@
services:
mongodb:
- image: mongo
+ image: mongo:8
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
diff --git a/global.json b/global.json
index 9935065..930c900 100644
--- a/global.json
+++ b/global.json
@@ -4,6 +4,6 @@
"rollForward": "latestFeature"
},
"msbuild-sdks": {
- "Cake.Sdk": "6.1.1"
+ "Cake.Sdk": "6.2.0"
}
}
diff --git a/src/MongoDb/Polling/ConfigurationImplementation.cs b/src/MongoDb/Polling/ConfigurationImplementation.cs
index 1c496db..a272851 100644
--- a/src/MongoDb/Polling/ConfigurationImplementation.cs
+++ b/src/MongoDb/Polling/ConfigurationImplementation.cs
@@ -161,12 +161,15 @@ public void ConfigureServices(OutboxKey key, IServiceCollection services)
services.AddKeyedSingleton(key, _dbFactory);
+ services.AddKeyedSingleton(key);
+
services.AddKeyedSingleton(
key,
(s, _) => ActivatorUtilities.CreateInstance(
s,
_lockBaseSettings,
- s.GetRequiredKeyedService>(key)(key, s)));
+ s.GetRequiredKeyedService>(key)(key, s),
+ s.GetRequiredKeyedService(key)));
_collectionConfigurator.ConfigureMe(new GetMongoDbOutboxCollectionConfigured(
key,
diff --git a/src/MongoDb/Synchronization/ChangeStreamListener.cs b/src/MongoDb/Synchronization/ChangeStreamListener.cs
index 3a3951e..c6827a1 100644
--- a/src/MongoDb/Synchronization/ChangeStreamListener.cs
+++ b/src/MongoDb/Synchronization/ChangeStreamListener.cs
@@ -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> cursor,
- CancellationTokenSource cts) : IAsyncDisposable
+internal interface IChangeStreamNotifier : IAsyncDisposable
{
- public Task WaitAsync() => autoResetEvent.WaitAsync(cts.Token);
+ Task OnChangeAsync(CancellationToken ct);
+}
- public static async Task StartAsync(
+internal sealed partial class ChangeStreamListener(ILogger logger)
+{
+ public async Task ListenAsync(
IMongoCollection collection,
DistributedLockDefinition lockDefinition,
CancellationToken ct)
@@ -21,41 +22,140 @@ public static async Task 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 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> 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);
}
\ No newline at end of file
diff --git a/src/MongoDb/Synchronization/DistributedLockThingy.cs b/src/MongoDb/Synchronization/DistributedLockThingy.cs
index 65aa66b..902273b 100644
--- a/src/MongoDb/Synchronization/DistributedLockThingy.cs
+++ b/src/MongoDb/Synchronization/DistributedLockThingy.cs
@@ -7,6 +7,7 @@ internal sealed partial class DistributedLockThingy(
DistributedLockSettings settings,
IMongoDatabase database,
TimeProvider timeProvider,
+ ChangeStreamListener changeStreamListener,
ILogger logger)
{
private readonly IMongoCollection _collection =
@@ -28,7 +29,7 @@ public async Task AcquireAsync(DistributedLockDefinition lockD
}
catch (Exception)
{
- await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
+ await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
throw;
}
}
@@ -42,7 +43,7 @@ public async Task AcquireAsync(DistributedLockDefinition lockD
{
if (!await InnerTryAcquireAsync(internalLockDefinition, ct))
{
- await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
+ await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
return null;
}
@@ -52,7 +53,7 @@ public async Task AcquireAsync(DistributedLockDefinition lockD
}
catch (Exception)
{
- await internalLockDefinition.ChangeStreamListener.TryDisposeAsync();
+ await internalLockDefinition.ChangeStreamNotifier.TryDisposeAsync();
throw;
}
}
@@ -61,14 +62,14 @@ private async Task 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
};
}
@@ -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;
}
}
@@ -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)
{
@@ -333,6 +336,13 @@ private sealed class DistributedLock(
CancellationTokenSource keepAliveCts,
Func 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();
+ }
+ }
}
}
\ No newline at end of file
diff --git a/src/MongoDb/Synchronization/InternalDistributedLockDefinition.cs b/src/MongoDb/Synchronization/InternalDistributedLockDefinition.cs
index 9a774b5..bc7acb7 100644
--- a/src/MongoDb/Synchronization/InternalDistributedLockDefinition.cs
+++ b/src/MongoDb/Synchronization/InternalDistributedLockDefinition.cs
@@ -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;
diff --git a/tests/MongoDb.Tests/MongoDb.Tests.csproj b/tests/MongoDb.Tests/MongoDb.Tests.csproj
index f9e94df..8a60c2f 100644
--- a/tests/MongoDb.Tests/MongoDb.Tests.csproj
+++ b/tests/MongoDb.Tests/MongoDb.Tests.csproj
@@ -21,6 +21,7 @@
runtime; build; native; contentfiles; analyzers; buildtransitive
+
diff --git a/tests/MongoDb.Tests/Polling/BatchFetcherTests.cs b/tests/MongoDb.Tests/Polling/BatchFetcherTests.cs
index 7d72b8a..5136580 100644
--- a/tests/MongoDb.Tests/Polling/BatchFetcherTests.cs
+++ b/tests/MongoDb.Tests/Polling/BatchFetcherTests.cs
@@ -167,6 +167,7 @@ private IBatchFetcher CreateSut(
},
_db,
timeProvider,
+ new ChangeStreamListener(NullLogger.Instance),
NullLogger.Instance),
new BatchCompleter(Defaults.Delete.MongoDbPollingSettings,
Defaults.Delete.CollectionConfig,
@@ -188,6 +189,7 @@ private IBatchFetcher CreateSut(
},
_db,
timeProvider,
+ new ChangeStreamListener(NullLogger.Instance),
NullLogger.Instance),
new BatchCompleter(Defaults.Update.MongoDbPollingSettings,
Defaults.Update.CollectionConfigWithProcessedAt,
diff --git a/tests/MongoDb.Tests/Synchronization/DistributedLockThingyTests.cs b/tests/MongoDb.Tests/Synchronization/DistributedLockThingyTests.cs
index f33c68f..7ff13c1 100644
--- a/tests/MongoDb.Tests/Synchronization/DistributedLockThingyTests.cs
+++ b/tests/MongoDb.Tests/Synchronization/DistributedLockThingyTests.cs
@@ -1,5 +1,6 @@
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
+using Microsoft.Extensions.Logging.Testing;
using Microsoft.Extensions.Time.Testing;
using MongoDB.Driver;
using YakShaveFx.OutboxKit.MongoDb.Synchronization;
@@ -9,6 +10,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Tests.Synchronization;
public class DistributedLockThingyTests(MongoDbFixture fixture)
{
private readonly ILogger _logger = NullLogger.Instance;
+ private readonly ILoggerFactory _loggerFactory = new NullLoggerFactory();
private readonly string _databaseName = $"test_{Guid.NewGuid():N}";
private readonly CancellationToken _ct = TestContext.Current.CancellationToken;
@@ -18,7 +20,7 @@ public async Task WhenAcquiringAvailableLockThenItsAcquired()
var database = GetDatabase();
var settings = new DistributedLockSettings { ChangeStreamsEnabled = false };
var timeProvider = new FakeTimeProvider();
- var sut = new DistributedLockThingy(settings, database, timeProvider, _logger);
+ var sut = new DistributedLockThingy(settings, database, timeProvider, GetChangeStreamListener(), _logger);
var lockDef = CreateLockDefinition();
await using var @lock = await sut.AcquireAsync(lockDef, _ct);
@@ -36,8 +38,8 @@ public async Task WhenAcquiringUnavailableLockThenItWaitsForExpiration()
// so we can control time separately and test expiration
var sut1TimeProvider = new FakeTimeProvider();
var sut2TimeProvider = new FakeTimeProvider();
- var sut1 = new DistributedLockThingy(settings, GetDatabase(), sut1TimeProvider, _logger);
- var sut2 = new DistributedLockThingy(settings, GetDatabase(), sut2TimeProvider, _logger);
+ var sut1 = new DistributedLockThingy(settings, GetDatabase(), sut1TimeProvider, GetChangeStreamListener(), _logger);
+ var sut2 = new DistributedLockThingy(settings, GetDatabase(), sut2TimeProvider, GetChangeStreamListener(), _logger);
var lockDef = CreateLockDefinition();
// acquire first lock
@@ -63,8 +65,8 @@ public async Task WhenAcquiringUnavailableLockThenItWaitsForChangeStreamsNotific
{
var settings = new DistributedLockSettings { ChangeStreamsEnabled = true };
var timeProvider = new FakeTimeProvider();
- var sut1 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
- var sut2 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
+ var sut1 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
+ var sut2 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
var lockDef = CreateLockDefinition();
var acquireLock1Tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
var releaseLock1Tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
@@ -100,8 +102,8 @@ public async Task WhenTryAcquireUnavailableLockThenItReturnsNull()
{
var settings = new DistributedLockSettings { ChangeStreamsEnabled = false };
var timeProvider = new FakeTimeProvider();
- var sut1 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
- var sut2 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
+ var sut1 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
+ var sut2 = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
var lockDef = CreateLockDefinition();
// acquire first lock
@@ -130,7 +132,7 @@ public async Task WhenTryAcquireWithConcurrentAttemptsThenOnlyOneSucceeds()
var acquiredLocks = await Task.WhenAll(
attempts.Select(def =>
{
- var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
+ var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
return sut.TryAcquireAsync(def, _ct);
}));
@@ -144,8 +146,8 @@ public async Task WhenLockIsLostThenItsDetectedOnNextKeepAlive()
// so we can control time separately and test expiration
var sut1TimeProvider = new FakeTimeProvider();
var sut2TimeProvider = new FakeTimeProvider();
- var sut1 = new DistributedLockThingy(settings, GetDatabase(), sut1TimeProvider, _logger);
- var sut2 = new DistributedLockThingy(settings, GetDatabase(), sut2TimeProvider, _logger);
+ var sut1 = new DistributedLockThingy(settings, GetDatabase(), sut1TimeProvider, GetChangeStreamListener(), _logger);
+ var sut2 = new DistributedLockThingy(settings, GetDatabase(), sut2TimeProvider, GetChangeStreamListener(), _logger);
var lockLostCalled = false;
var lockDef = CreateLockDefinition() with
{
@@ -181,7 +183,7 @@ public async Task WhenLockIsLostThenItsDetectedByChangeStreams()
{
var settings = new DistributedLockSettings { ChangeStreamsEnabled = true };
var timeProvider = new FakeTimeProvider();
- var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
+ var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
var lockLostTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
var lock1AcquiredTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
var lockDef = CreateLockDefinition() with
@@ -216,7 +218,7 @@ public async Task WhenLockIsDeletedExternallyThenItsDetectedByChangeStreamsAndRe
{
var settings = new DistributedLockSettings { ChangeStreamsEnabled = true };
var timeProvider = new FakeTimeProvider();
- var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, _logger);
+ var sut = new DistributedLockThingy(settings, GetDatabase(), timeProvider, GetChangeStreamListener(), _logger);
var lockLostCalled = false;
var lockDef = CreateLockDefinition() with
{
@@ -258,7 +260,7 @@ public async Task WhenAcquiringAndReleasingLockThenItsReleased()
var database = GetDatabase();
var settings = new DistributedLockSettings { ChangeStreamsEnabled = false };
var timeProvider = TimeProvider.System;
- var sut = new DistributedLockThingy(settings, database, timeProvider, _logger);
+ var sut = new DistributedLockThingy(settings, database, timeProvider, GetChangeStreamListener(), _logger);
var lockDef = CreateLockDefinition() with { Duration = TimeSpan.FromSeconds(1) };
var @lock = await sut.AcquireAsync(lockDef, _ct);
@@ -277,7 +279,12 @@ public async Task WhenAcquiringAndReleasingLockWithChangeStreamsEnabledThenItsRe
var database = GetDatabase();
var settings = new DistributedLockSettings { ChangeStreamsEnabled = true };
var timeProvider = new FakeTimeProvider();
- var sut = new DistributedLockThingy(settings, database, timeProvider, _logger);
+ var sut = new DistributedLockThingy(
+ settings,
+ database,
+ timeProvider,
+ new ChangeStreamListener(_loggerFactory.CreateLogger()),
+ _logger);
var lockDef = CreateLockDefinition();
var @lock = await sut.AcquireAsync(lockDef, _ct);
@@ -292,7 +299,31 @@ public async Task WhenAcquiringAndReleasingLockWithChangeStreamsEnabledThenItsRe
doc.Should().BeNull();
}
- // need to delete using a different client, otherwise it won't trigger change streams
+ [Fact]
+ public async Task WhenKeepAliveRenewsLockWithChangeStreamsThenPotentiallyLostIsNotLogged()
+ {
+ var settings = new DistributedLockSettings { ChangeStreamsEnabled = true };
+ var fakeLogCollector = new FakeLogCollector();
+ var fakeLogger = new FakeLogger(fakeLogCollector);
+ var sut = new DistributedLockThingy(
+ settings,
+ GetDatabase(),
+ TimeProvider.System,
+ new ChangeStreamListener(fakeLogger),
+ _logger);
+ var lockDef = CreateLockDefinition() with { Duration = TimeSpan.FromSeconds(2) };
+
+ await using var @lock = await sut.AcquireAsync(lockDef, _ct);
+
+ // let keepalive fire a few times (self-replace events hit change stream)
+ await Task.Delay(TimeSpan.FromSeconds(5), _ct);
+
+ fakeLogCollector
+ .GetSnapshot()
+ .Should()
+ .NotContain(r => r.Message.Contains("potentially lost"));
+ }
+
private async Task DeleteLock(DistributedLockDefinition lockDef)
=> await GetDatabase()
.GetCollection()
@@ -332,6 +363,8 @@ private async Task ReplaceOwner(DistributedLockDefinition lock
private IMongoDatabase GetDatabase()
=> new MongoClient(fixture.ConnectionString).GetDatabase(_databaseName);
+
+ private ChangeStreamListener GetChangeStreamListener() => new (_loggerFactory.CreateLogger());
}
file static class Extensions