BG.Common.Inbox
1.0.0
dotnet add package BG.Common.Inbox --version 1.0.0
NuGet\Install-Package BG.Common.Inbox -Version 1.0.0
<PackageReference Include="BG.Common.Inbox" Version="1.0.0" />
<PackageVersion Include="BG.Common.Inbox" Version="1.0.0" />
<PackageReference Include="BG.Common.Inbox" />
paket add BG.Common.Inbox --version 1.0.0
#r "nuget: BG.Common.Inbox, 1.0.0"
#:package BG.Common.Inbox@1.0.0
#addin nuget:?package=BG.Common.Inbox&version=1.0.0
#tool nuget:?package=BG.Common.Inbox&version=1.0.0
BG.Common.Inbox
Generyczny, multi-instance wzorzec Inbox. Pozwala dowolnemu modułowi/serwisowi
zarejestrować własny, w pełni odizolowany InboxProcessor (własny store, własne
strategie przetwarzania, własna konfiguracja) w tym samym hoście DI — analogicznie
do BG.Common.Outbox, ale dla wiadomości przychodzących.
Co dostarcza pakiet
InboxMessage— model odebranej wiadomości.Idto identyfikator nadany przez system źródłowy (np. id wiadomości z brokera), a nie generowany lokalnie — to on jest kluczem deduplikacji.IInboxStore— store/repozytorium wystawione jako pierwszoklasowa abstrakcja, nie tylko szczegół implementacyjny procesora. Konsument może wstrzyknąć ją bezpośrednio (np. w konsumerze brokera) i wywołać:ExistsAsync(id)— szybki, tani check czy dana wiadomość już była widziana,TryAddAsync(message)— atomowy, idempotentny zapis; zwracafalse(bez wyjątku) gdy wiadomość o tymIdjuż istnieje, co jest głównym mechanizmem ochrony przed podwójnym przetworzeniem tej samej wiadomości.- Store dostarcza też
ClaimBatchAsync/MarkAsProcessedAsync/MarkAsFailedAsync/ReleaseLocksAsync/PurgeProcessedAsyncna potrzebyInboxProcessor.
IInboxMessageProcessingStrategy— strategia przetwarzania jednego rodzaju wiadomości (CanProcess+ProcessAsync).InboxProcessorwybiera pierwszą pasującą zarejestrowaną strategię dla danej wiadomości — jedna implementacja naMessageType.InboxProcessorOptions— konfiguracja (BatchSize,PollingInterval,LockDuration,MaxAttempts,RetentionPeriod,PurgeInterval).InboxProcessor—BackgroundServiceodpytujący store, dobierający strategię i przetwarzający wiadomości.AddInboxProcessor(...)— rejestracja jednej, nazwanej/kluczowanej instancji procesora.
Podobnie jak w Outboksie, store, strategie i opcje są kluczowane/nazwane tym samym
key, więc wiele instancji InboxProcessor może współistnieć w jednym
IServiceProvider bez wzajemnego mieszania danych ani konfiguracji.
Użycie — pojedynczy moduł
services.AddKeyedScoped<IInboxStore, MyInboxStore>("orders");
services.AddKeyedScoped<IInboxMessageProcessingStrategy, OrderCreatedStrategy>("orders");
services.AddKeyedScoped<IInboxMessageProcessingStrategy, OrderCancelledStrategy>("orders");
services.AddInboxProcessor("orders", configuration, options =>
{
options.BatchSize = 50;
});
Deduplikacja w miejscu odbioru wiadomości (np. w konsumerze brokera), zanim
InboxProcessor w ogóle zdąży ją przetworzyć:
public class OrderEventsConsumer(
[FromKeyedServices("orders")] IInboxStore inboxStore)
{
public async Task OnMessageAsync(BrokerMessage brokerMessage, CancellationToken ct)
{
var message = InboxMessage.Create(brokerMessage.MessageId, brokerMessage.Type, brokerMessage.Body);
if (!await inboxStore.TryAddAsync(message, ct))
{
// Już widzieliśmy tę wiadomość — pomijamy, ale ack-ujemy broker.
return;
}
// InboxProcessor podejmie zapisaną wiadomość w kolejnym cyklu i przetworzy
// ją właściwą IInboxMessageProcessingStrategy.
}
}
Konfiguracja w appsettings.json (opcjonalnie, configure nadpisuje wartości z sekcji):
{
"Inbox": {
"Processor": {
"orders": {
"BatchSize": 50,
"PollingInterval": "00:00:05",
"LockDuration": "00:00:30",
"MaxAttempts": 5
}
}
}
}
Użycie — dwa niezależne moduły w jednym hoście
// Moduł "orders"
services.AddKeyedScoped<IInboxStore, OrdersInboxStore>("orders");
services.AddKeyedScoped<IInboxMessageProcessingStrategy, OrderCreatedStrategy>("orders");
services.AddInboxProcessor("orders", configuration, o => o.BatchSize = 50);
// Moduł "notifications"
services.AddKeyedScoped<IInboxStore, NotificationsInboxStore>("notifications");
services.AddKeyedScoped<IInboxMessageProcessingStrategy, NotificationReceivedStrategy>("notifications");
services.AddInboxProcessor("notifications", configuration, o => o.BatchSize = 10);
orders i notifications działają jako dwie niezależne instancje InboxProcessor
w tym samym IServiceProvider — każda korzysta wyłącznie ze swojego store'a, swoich
strategii przetwarzania i swojej konfiguracji.
Czyszczenie przetworzonych wiadomości
W przeciwieństwie do Outboksa,
gdzie po opublikowaniu zdarzenia wiersz można od razu skasować, Inbox musi
trzymać przetworzone wiadomości przez jakiś czas — to właśnie one chronią przed
ponownym przetworzeniem tej samej wiadomości dostarczonej powtórnie (at-least-once
delivery). Natychmiastowe usuwanie wiersza w MarkAsProcessedAsync zepsułoby
deduplikację: redelivered duplikat już przetworzonej wiadomości nie zostałby
wykryty przez TryAddAsync/ExistsAsync i zostałby przetworzony ponownie.
Dlatego InboxProcessor czyści store okresowo, a nie natychmiast po przetworzeniu:
InboxProcessorOptions.RetentionPeriod(domyślnie 7 dni) — jak długo przetworzona wiadomość musi być trzymana, zanim stanie się kandydatem do usunięcia. Dobierz tę wartość do realistycznego okna redelivery Twojego brokera/transportu.InboxProcessorOptions.PurgeInterval(domyślnie 1 godzina) — jak często procesor wołaIInboxStore.PurgeProcessedAsync(DateTimeOffset.UtcNow - RetentionPeriod), niezależnie odPollingInterval(żeby nie odpytywać o czyszczenie przy każdym cyklu).
services.AddInboxProcessor("orders", configuration, options =>
{
options.RetentionPeriod = TimeSpan.FromDays(3);
options.PurgeInterval = TimeSpan.FromHours(1);
});
Implementacja IInboxStore.PurgeProcessedAsync powinna usuwać wiersze przetworzone
przed przekazaną granicą czasową (np. DELETE ... WHERE ProcessedAt < @olderThan),
utrzymując tabelę w rozsądnym rozmiarze bez rezygnowania z deduplikacji w realistycznym oknie.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Options (>= 10.0.0)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 1.0.0 | 129 | 7/15/2026 |