Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
9479991
spec: 017 phase 4 task 4.1 — phase 4 predicted; a verdict per hard block
iancooper Sep 27, 2026
bd2b1ed
blockcheck: grow the pin by the Npgsql EF provider and two TickerQ pa…
iancooper Sep 27, 2026
fc65dce
blockcheck: baseline TickerQScheduler.md #3, built by the grown pin
iancooper Sep 27, 2026
a8ede7c
docs: outbox and inbox pages compile; the InMemory Inbox's window and…
iancooper Sep 27, 2026
f8c21aa
blockcheck: baseline the twenty-four blocks 4.2 made build
iancooper Sep 27, 2026
7d04918
spec: 017 phase 4 task 4.2 — outbox and inbox pages whole; #4335 stated
iancooper Sep 27, 2026
b8f21ac
docs: the six distributed-lock pages compile; session release run on …
iancooper Sep 27, 2026
05b8916
blockcheck: baseline the twelve blocks 4.3 made build
iancooper Sep 27, 2026
05fdeaf
spec: 017 phase 4 task 4.3 — lock pages whole; the ≤ 60 target met at 58
iancooper Sep 27, 2026
2defac6
docs: the sweeper circuit-breaking and Azure Blob archive pages compile
iancooper Sep 27, 2026
d47bfdf
blockcheck: baseline the five blocks 4.4 made build
iancooper Sep 27, 2026
cdeabca
spec: 017 phase 4 task 4.4 — the remaining outbox pages whole; 57 wit…
iancooper Sep 27, 2026
28b2d5f
docs: configure the sweeper's interval as 10.7.0 does, and count the …
iancooper Sep 27, 2026
6263986
blockcheck: baseline the two sweeper-interval blocks
iancooper Sep 27, 2026
e7cec73
spec: 017 phase 4 task 4.4 — SweeperCircuitBreaking.md repaired by ru…
iancooper Sep 27, 2026
a61893b
docs: say which Outboxes honour a tripped topic, and what explicit cl…
iancooper Sep 27, 2026
2a24787
blockcheck: baseline SweeperCircuitBreaking.md #5, and re-admit the p…
iancooper Sep 27, 2026
9e129dd
spec: 017 phase 4 task 4.4 — SweeperCircuitBreaking.md's outbox and c…
iancooper Sep 27, 2026
be7793c
spec: 017 task 4.4 — the two ledger rows' before-counts, measured apart
iancooper Sep 27, 2026
ca6a0b2
docs: link the DynamoDB and Spanner tripped-topic issues, #4443 and #…
iancooper Sep 27, 2026
b3fcc64
spec: 017 phase 4 task 4.5 — phase 4 closed; rows 2 and 9 moved
iancooper Sep 28, 2026
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
96 changes: 68 additions & 28 deletions contents/AzureBlobArchiveProvider.md
Original file line number Diff line number Diff line change
@@ -1,49 +1,89 @@
---
description: "The Azure Blob Archive Provider is a provider for Outbox Archiver."
description: "The Azure Blob Archive Provider writes the messages the Outbox Archiver takes from your Outbox into an Azure Blob Storage container."
layout:
description:
visible: false
---

# Azure Blob Archive Provider

> **Reference** · Applies to **Brighter V10**
> **Reference** · Applies to **Brighter V10** · Prerequisites: [Outbox Archiver](/contents/OutboxArchiver.md)

## Azure Blob Archive Provider Usage
The Azure Blob Archive Provider is a provider for [Outbox Archiver](/contents/OutboxArchiver.md).
The Azure Blob Archive Provider writes the messages the Outbox Archiver takes from your Outbox into an Azure Blob Storage container. It is an `IAmAnArchiveProvider` you pass to `UseOutboxArchiver<TTransaction>`; the [Outbox Archiver](/contents/OutboxArchiver.md) decides when a message is old enough to archive, and removes it from the Outbox once the provider has written it.

For this we will need the *Archive* packages for the Azure *Archive Provider*.
## Azure Blob Archive Provider Configuration

* **Paramore.Brighter.Archive.Azure**
You need three packages:

* **Paramore.Brighter.Archive.Azure** — the provider, in the `Paramore.Brighter.Storage.Azure` namespace
* **Paramore.Brighter.Outbox.Hosting** — `UseOutboxArchiver`, which runs the Archiver as a hosted service
* **Azure.Identity** — a `TokenCredential` for the provider to write with. The provider package does not bring it in

```csharp
using System;
using System.Data.Common;
using Azure.Identity;
using Azure.Storage.Blobs.Models;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Paramore.Brighter.Extensions.DependencyInjection;
using Paramore.Brighter.Outbox.Hosting;
using Paramore.Brighter.Storage.Azure;

``` csharp
private static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices(hostContext, services) =>
.ConfigureServices((hostContext, services) =>
{
ConfigureBrighter(hostContext, services);
}
ConfigureBrighter(services);
});

private static void ConfigureBrighter(HostBuilderContext hostContext, IServiceCollection services)
private static void ConfigureBrighter(IServiceCollection services)
{
services.AddBrighter(options =>
{ ... })
.UseOutboxArchiver(
new AzureBlobArchiveProvider(new AzureBlobArchiveProviderOptions()
{
BlobContainerUri = "https://brighterarchivertest.blob.core.windows.net/messagearchive",
TokenCredential = New AzCliCredential();
}
),
options => {
TimerInterval = 5; // Every 5 seconds
BatchSize = 500; // 500 messages at a time
MinimumAge = 744; // 1 month
}
);
services.AddBrighter()
.AddProducers(configure =>
{
// ... your producer registry, and the Outbox the Archiver reads from
})
// DbTransaction is the transaction type of a relational Outbox; see Outbox Archiver for the others
.UseOutboxArchiver<DbTransaction>(
new AzureBlobArchiveProvider(new AzureBlobArchiveProviderOptions(
blobContainerUri: new Uri("https://brighterarchivertest.blob.core.windows.net/messagearchive"),
tokenCredential: new AzureCliCredential(),
accessTier: AccessTier.Cool,
tagBlobs: true)),
options =>
{
options.TimerInterval = 5; // every 5 seconds
options.ArchiveBatchSize = 500; // 500 messages at a time
options.MinimumAge = TimeSpan.FromDays(31); // dispatched more than a month ago
});
}
```

...
`TTransaction` is your Outbox's transaction type, not its transaction provider; [Outbox Archiver](/contents/OutboxArchiver.md) lists the type for each Outbox, and the options the second argument sets.

```
## Azure Blob Archive Provider Options

`AzureBlobArchiveProviderOptions` takes its first four values as constructor arguments; the rest have defaults. Its properties are `init`-only, so set them, and the two functions, in an object initializer when you create the options.

| Option | Type | Default | Description |
|---|---|---|---|
| `BlobContainerUri` | `Uri` | required | The container the provider writes to. The provider does not create it |
| `TokenCredential` | `TokenCredential` | required | The credential the provider writes with — `AzureCliCredential` locally, a managed identity or `DefaultAzureCredential` in Azure |
| `AccessTier` | `AccessTier` | required | The access tier each blob is written in, such as `Hot`, `Cool` or `Archive` |
| `TagBlobs` | `bool` | required | Whether to write index tags on each blob, from `TagsFunc` |
| `MaxConcurrentUploads` | `int` | `8` | The most transfers one upload runs in parallel |
| `MaxUploadSize` | `int` | `50` | The largest chunk one transfer sends, in megabytes |
| `TagsFunc` | `Func<Message, Dictionary<string, string?>>` | the message's topic, correlation id, message type, timestamp and content type | The tags written when `TagBlobs` is `true` |
| `StorageLocationFunc` | `Func<Message, string>` | the message's Id | The name of the blob a message is written to, within the container |

## What the Azure Blob Archive Provider Writes

Each archived message becomes one blob, named by `StorageLocationFunc` — by default, the message's Id. The blob holds the message **body**; the header reaches the archive only through the tags, and only when `TagBlobs` is `true`. If you need more of the header, write it into the tags with your own `TagsFunc`.

A message whose blob already exists is not written again, so archiving a message twice is harmless.

## Further Reading

- [Outbox Archiver](/contents/OutboxArchiver.md) - When messages are archived, and the `TTransaction` for each Outbox
- [Outbox Support](/contents/BrighterOutboxSupport.md) - The Outbox, the Sweeper and the Archiver
14 changes: 13 additions & 1 deletion contents/AzureBlobDistributedLock.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@ Configure the provider with `AzureBlobLockingProvider`, passing an
`TokenCredential`:

```csharp
using System;
using Azure.Identity;
using Paramore.Brighter.Locking.Azure;

new AzureBlobLockingProvider(
new AzureBlobLockingProviderOptions(
blobContainerUri: new Uri("https://myaccount.blob.core.windows.net/brighter-locks"),
Expand Down Expand Up @@ -56,11 +60,19 @@ initialiser like any other member.
## Azure Blob Distributed Lock Example

```csharp
using System;
using Azure.Identity;
using Microsoft.Extensions.DependencyInjection;
using Paramore.Brighter;
using Paramore.Brighter.Extensions.DependencyInjection;
using Paramore.Brighter.Locking.Azure;
using Paramore.Brighter.Outbox.Hosting;

services
.AddBrighter()
.AddProducers(opt =>
{
opt.Outbox = /* your external Outbox */;
opt.Outbox = outbox; // your external Outbox
// ... connection/transaction providers for your Outbox ...

opt.DistributedLock = new AzureBlobLockingProvider(
Expand Down
25 changes: 17 additions & 8 deletions contents/BrighterBasicConfiguration.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,39 +136,48 @@ We provide the class `ServiceActivatorHostedService` for this in the NuGet packa

The `ServiceActivatorHostedService` calls the **Dispatcher.Receive** method which starts message pumps for the configured *Subscriptions*.

``` csharp
```csharp
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Paramore.Brighter.ServiceActivator.Extensions.Hosting;

private static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices(hostContext, services) =>
.ConfigureServices((hostContext, services) =>
{
ConfigureBrighter(hostContext, services);
}
});

private static void ConfigureBrighter(HostBuilderContext hostContext, IServiceCollection services)
{
...
// ...
services.AddHostedService<ServiceActivatorHostedService>();
}

```

On shutdown Brighter will allow the current *Request Handler* to complete, then end the message pump loop and exit. If you have long-running handlers it is possible that they will not complete in the default 5s for graceful shutdown of the MS Generic Host. In this case, you need to [increase the timeout](https://docs.microsoft.com/en-us/aspnet/core/fundamentals/host/generic-host?view=aspnetcore-6.0#shutdowntimeout) of the host shutdown.

``` csharp
```csharp
using System;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Paramore.Brighter.ServiceActivator.Extensions.Hosting;

private static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices(hostContext, services) =>
.ConfigureServices((hostContext, services) =>
{
services.Configure<HostOptions>(options =>
{
options.ShutdownTimeout = TimeSpan.FromSeconds(20);
});
ConfigureBrighter(hostContext, services);
}
});

private static void ConfigureBrighter(HostBuilderContext hostContext, IServiceCollection services)
{
...
// ...

services.AddHostedService<ServiceActivatorHostedService>();
}
Expand Down
6 changes: 6 additions & 0 deletions contents/BrighterInboxSupport.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,12 @@ There are two versions of the attribute: sync and async. Ensure that you choose

Your inbox is configured as part of the Brighter extensions to `IServiceCollection`. See [Inbox Configuration](/contents/DispatcherConfigurationReference.md#inbox) for more.

### Global Inbox Configuration in a Consumer-Only Application

In Brighter 10.7.0, the `InboxConfiguration` you set in `AddConsumers` reaches your handlers' pipelines only when the application also calls `AddProducers`. In an application that only consumes, Brighter adds no Inbox to any handler, so a duplicate is handled again, with no error. This is [BrighterCommand/Brighter#4335](https://github.com/BrighterCommand/Brighter/issues/4335); the fix is on Brighter's development branch and in no release yet.

Until a release carries it, put the attribute on each handler that must not see a duplicate, as in [Adding an Inbox to a Handler](#adding-an-inbox-to-a-handler). The attribute takes effect whether or not the application registers producers. Its own `onceOnlyAction` decides what a duplicate does, whatever `actionOnExists` the configuration sets, so give it the action you want. An application that calls `AddProducers` gets the global Inbox as configured.

### Provisioning the Inbox Table

If your Inbox runs on a relational database (MSSQL, PostgreSQL, MySQL, SQLite, or Spanner), Brighter can create and migrate the table for you at application startup — see [Database Provisioning](/contents/BoxProvisioning.md). The **Inbox Builder** section below describes the alternative: managing the DDL yourself.
Expand Down
104 changes: 61 additions & 43 deletions contents/DapperOutbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ public void ConfigureServices(IServiceCollection services)

```

The handler below is `AddGreetingHandlerAsync` from the Brighter sample at `Brighter/samples/WebAPI/WebAPI_Dapper/`, which also declares the `AddGreeting` request, the `Person` and `Greeting` entities and the `GreetingMade` event it uses.

In our handler we take a dependency on Brighter's **IAmATransactionConnectionProvider**. We explicitly start a transaction within the handler on the Database within the provider. Dapper provides extension methods on a DbConnection for typical CRUD operations. Our provider wraps that DbConnection, and allows you to create a DB transaction associated with that DbConnection. You must use our method, and not create the transaction directly via the connection, because we cannot obtain that transaction. Sharing that transaction allows us to insert a message into the Outbox within the same transaction.

We call **DepositPostAsync** within that transaction to write the message to the Outbox. Once the transaction has closed we can call **ClearOutboxAsync** to immediately clear, or we can rely on the Outbox Sweeper, if we have configured one to clear for us. (There are equivalent synchronous versions of these APIs).
Expand All @@ -78,56 +80,72 @@ using Dapper;
using Microsoft.Extensions.Logging;
using Paramore.Brighter;

public override async Task<AddGreeting> HandleAsync(AddGreeting addGreeting, CancellationToken cancellationToken = default)
public class AddGreetingHandlerAsync : RequestHandlerAsync<AddGreeting>
{
var posts = new List<Id>();
private readonly IAmATransactionConnectionProvider _transactionProvider;
private readonly IAmACommandProcessor _postBox;
private readonly ILogger<AddGreetingHandlerAsync> _logger;

//We use the transaction provider to grab connection and transaction, because Outbox needs
//to share them 'behind the scenes'
DbConnection conn = await _transactionProvider.GetConnectionAsync(cancellationToken);
DbTransaction tx = await _transactionProvider.GetTransactionAsync(cancellationToken);
try
{
var people = await conn.QueryAsync<Person>(
"select * from Person where name = @name",
new { name = addGreeting.Name },
tx);
var person = people.Single();

var greeting = new Greeting(addGreeting.Greeting, person);

//write the added child entity to the Db
await conn.ExecuteAsync(
"insert into Greeting (Message, Recipient_Id) values (@Message, @RecipientId)",
new { greeting.Message, greeting.RecipientId },
tx);

//Now write the message we want to send to the Db in the same transaction.
posts.Add(await _postBox.DepositPostAsync(
new GreetingMade(greeting.Greet()),
_transactionProvider,
cancellationToken: cancellationToken));

//commit both new greeting and outgoing message
await _transactionProvider.CommitAsync(cancellationToken);
}
catch (Exception e)
public AddGreetingHandlerAsync(IAmATransactionConnectionProvider transactionProvider,
IAmACommandProcessor postBox,
ILogger<AddGreetingHandlerAsync> logger)
{
_logger.LogError(e, "Exception thrown handling Add Greeting request");
//it went wrong, rollback the entity change and the downstream message
await _transactionProvider.RollbackAsync(cancellationToken);
return await base.HandleAsync(addGreeting, cancellationToken);
_transactionProvider = transactionProvider;
_postBox = postBox;
_logger = logger;
}
finally

public override async Task<AddGreeting> HandleAsync(AddGreeting addGreeting, CancellationToken cancellationToken = default)
{
_transactionProvider.Close();
}
var posts = new List<Id>();

//Send this message via a transport. We need the ids to send just the messages here, not all outstanding ones.
//Alternatively, you can let the Sweeper do this, but at the cost of increased latency
await _postBox.ClearOutboxAsync(posts, cancellationToken: cancellationToken);
//We use the transaction provider to grab connection and transaction, because Outbox needs
//to share them 'behind the scenes'
DbConnection conn = await _transactionProvider.GetConnectionAsync(cancellationToken);
DbTransaction tx = await _transactionProvider.GetTransactionAsync(cancellationToken);
try
{
var people = await conn.QueryAsync<Person>(
"select * from Person where name = @name",
new { name = addGreeting.Name },
tx);
var person = people.Single();

var greeting = new Greeting(addGreeting.Greeting, person);

//write the added child entity to the Db
await conn.ExecuteAsync(
"insert into Greeting (Message, Recipient_Id) values (@Message, @RecipientId)",
new { greeting.Message, greeting.RecipientId },
tx);

//Now write the message we want to send to the Db in the same transaction.
posts.Add(await _postBox.DepositPostAsync(
new GreetingMade(greeting.Greet()),
_transactionProvider,
cancellationToken: cancellationToken));

//commit both new greeting and outgoing message
await _transactionProvider.CommitAsync(cancellationToken);
}
catch (Exception e)
{
_logger.LogError(e, "Exception thrown handling Add Greeting request");
//it went wrong, rollback the entity change and the downstream message
await _transactionProvider.RollbackAsync(cancellationToken);
return await base.HandleAsync(addGreeting, cancellationToken);
}
finally
{
_transactionProvider.Close();
}

return await base.HandleAsync(addGreeting, cancellationToken);
//Send this message via a transport. We need the ids to send just the messages here, not all outstanding ones.
//Alternatively, you can let the Sweeper do this, but at the cost of increased latency
await _postBox.ClearOutboxAsync(posts, cancellationToken: cancellationToken);

return await base.HandleAsync(addGreeting, cancellationToken);
}
}
```

Expand Down
Loading
Loading