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
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="BG.Common.Inbox" Version="1.0.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="BG.Common.Inbox" Version="1.0.0" />
                    
Directory.Packages.props
<PackageReference Include="BG.Common.Inbox" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add BG.Common.Inbox --version 1.0.0
                    
#r "nuget: BG.Common.Inbox, 1.0.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package BG.Common.Inbox@1.0.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=BG.Common.Inbox&version=1.0.0
                    
Install as a Cake Addin
#tool nuget:?package=BG.Common.Inbox&version=1.0.0
                    
Install as a Cake Tool

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. Id to 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; zwraca false (bez wyjątku) gdy wiadomość o tym Id już istnieje, co jest głównym mechanizmem ochrony przed podwójnym przetworzeniem tej samej wiadomości.
    • Store dostarcza też ClaimBatchAsync/MarkAsProcessedAsync/MarkAsFailedAsync/ ReleaseLocksAsync/PurgeProcessedAsync na potrzeby InboxProcessor.
  • IInboxMessageProcessingStrategy — strategia przetwarzania jednego rodzaju wiadomości (CanProcess + ProcessAsync). InboxProcessor wybiera pierwszą pasującą zarejestrowaną strategię dla danej wiadomości — jedna implementacja na MessageType.
  • InboxProcessorOptions — konfiguracja (BatchSize, PollingInterval, LockDuration, MaxAttempts, RetentionPeriod, PurgeInterval).
  • InboxProcessor — BackgroundService odpytują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ła IInboxStore.PurgeProcessedAsync(DateTimeOffset.UtcNow - RetentionPeriod), niezależnie od PollingInterval (ż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 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

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