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
6 changes: 6 additions & 0 deletions contents/BrighterBasicConfiguration.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ using Paramore.Brighter;
using Paramore.Brighter.Extensions.DependencyInjection;
using Paramore.Brighter.MessagingGateway.RMQ.Async;
using Paramore.Brighter.MySql;
using Microsoft.EntityFrameworkCore;
using Paramore.Brighter.MySql.EntityFrameworkCore;
using Paramore.Brighter.Outbox.Hosting;
using Paramore.Brighter.Outbox.MySql;
Expand All @@ -70,6 +71,9 @@ public void ConfigureServices(IServiceCollection services)
var outboxConfiguration = new RelationalDatabaseConfiguration(DbConnectionString());
services.AddSingleton<IAmARelationalDatabaseConfiguration>(outboxConfiguration);

services.AddDbContext<GreetingsEntityGateway>(options =>
options.UseMySql(DbConnectionString(), ServerVersion.AutoDetect(DbConnectionString())));

services.AddBrighter(options =>
{
options.HandlerLifetime = ServiceLifetime.Scoped;
Expand Down Expand Up @@ -238,6 +242,8 @@ private static void ConfigureBrighter(HostBuilderContext hostContext, IServiceCo
var outboxConfiguration = new RelationalDatabaseConfiguration(
DbConnectionString(), outBoxTableName: "Outbox");

services.AddSingleton<IAmARelationalDatabaseConfiguration>(outboxConfiguration);

services.AddConsumers(options =>
{
options.Subscriptions = subscriptions;
Expand Down
3 changes: 1 addition & 2 deletions contents/BrighterSchedulerSupport.md
Original file line number Diff line number Diff line change
Expand Up @@ -251,9 +251,8 @@ var subscription = new Subscription<ProcessOrderCommand>(

**Transports without Native Delay (Kafka, AWS SNS, etc.):**

- Requires an external scheduler (Quartz, Hangfire, etc.)
- Message is scheduled via the configured scheduler
- Falls back to immediate requeue if no scheduler configured
- If you configure none, Brighter's registration uses the in-memory scheduler, which does not survive a restart; configure a durable scheduler (Quartz, Hangfire, etc.) for production

## Choosing a Scheduler

Expand Down
17 changes: 13 additions & 4 deletions contents/CommandProcessorConfigurationReference.md
Original file line number Diff line number Diff line change
Expand Up @@ -622,14 +622,23 @@ An Outbox has three pieces:
In this example, we want to use EF Core with an MySQL Outbox. See the documentation for Outboxes for specific configuration options.

``` csharp
// ...
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Paramore.Brighter;
using Paramore.Brighter.Extensions.DependencyInjection;
using Paramore.Brighter.MySql;
using Paramore.Brighter.MySql.EntityFrameworkCore;
using Paramore.Brighter.Outbox.MySql;

public void ConfigureServices(IServiceCollection services)
{

var outboxConfiguration = new RelationalDatabaseConfiguration(DbConnectionString());
services.AddSingleton<IAmARelationalDatabaseConfiguration>(outboxConfiguration);

services.AddBrighter(...)
services.AddDbContext<GreetingsEntityGateway>(options =>
options.UseMySql(DbConnectionString(), ServerVersion.AutoDetect(DbConnectionString())));

services.AddBrighter()
.AddProducers((configure) =>
{
configure.Outbox = new MySqlOutbox(outboxConfiguration);
Expand All @@ -638,7 +647,7 @@ public void ConfigureServices(IServiceCollection services)
})
.AutoFromAssemblies();

...
// ...
}

```
Expand Down
4 changes: 3 additions & 1 deletion contents/ConfiguringOpenTelemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,9 @@ Brighter writes to that source only through a tracer, an `IAmABrighterTracer`, r
container. `AddBrighter()` does not register one, so listening to the source is not enough:
`AddSource("paramore.brighter")` on its own records no Brighter span.
`AddBrighterInstrumentation()`, from the `Paramore.Brighter.Extensions.Diagnostics` package, does
both: it registers the tracer and adds the source.
both: it registers the tracer and adds the source. The Command Processor takes that tracer only when
it has an external bus, configured with `AddProducers`; with no producers, its requests record no
span at 10.7.0. This is reported as [#4510](https://github.com/BrighterCommand/Brighter/issues/4510).

Use it on the tracer provider that `AddOpenTelemetry()` builds, which shares your application's
container. A provider built with `Sdk.CreateTracerProviderBuilder()` keeps its own services, so the
Expand Down
2 changes: 2 additions & 0 deletions contents/DapperOutbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ public void ConfigureServices(IServiceCollection services)
outBoxTableName: "outbox_messages",
inboxTableName: "inbox_messages");

services.AddSingleton<IAmARelationalDatabaseConfiguration>(configuration);

services.AddConsumers(options =>
{
options.InboxConfiguration = new InboxConfiguration(new MySqlInbox(configuration));
Expand Down
2 changes: 2 additions & 0 deletions contents/DynamoOutbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,8 @@ using Paramore.Brighter.Outbox.Hosting;
public void ConfigureServices(IServiceCollection services)
{
// ... dynamoDb is your IAmazonDynamoDB client, producerRegistry your transport
services.AddSingleton<IAmazonDynamoDB>(dynamoDb);

services.AddBrighter()
.AddProducers(configure =>
{
Expand Down
38 changes: 25 additions & 13 deletions contents/EFCoreOutbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,26 +29,38 @@ Obviously, {DB} should match. In the example below we use MySql, so we would nee

* **Paramore.Brighter.MySql**

As described in [Command Processor Configuration Reference](/contents/CommandProcessorConfigurationReference.md#outbox-support), we configure Brighter to use an outbox with the Use{DB}Outbox method call.
As described in [Command Processor Configuration Reference](/contents/CommandProcessorConfigurationReference.md#outbox-support), we configure Brighter to use an outbox in the `AddProducers` method call, by setting its `Outbox`.

As we want to use EF Core, we also call: Use{DB}TransactionConnectionProvider so that we can share your transaction scope when persisting messages to the outbox.
As we want to use EF Core, we also set its `TransactionProvider` to the EF Core transaction provider for your `DbContext`, so that we can share your transaction when persisting messages to the outbox. Brighter creates both providers from your container, so register what they are built from: the database configuration and the `DbContext`.

```csharp
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Paramore.Brighter;
using Paramore.Brighter.Extensions.DependencyInjection;
using Paramore.Brighter.MySql;
using Paramore.Brighter.MySql.EntityFrameworkCore;
using Paramore.Brighter.Outbox.Hosting;
using Paramore.Brighter.Outbox.MySql;

``` csharp
public void ConfigureServices(IServiceCollection services)
{
services.AddBrighter(...)
var outboxConfiguration = new RelationalDatabaseConfiguration(DbConnectionString());
services.AddSingleton<IAmARelationalDatabaseConfiguration>(outboxConfiguration);
services.AddDbContext<GreetingsEntityGateway>(options =>
options.UseMySql(DbConnectionString(), ServerVersion.AutoDetect(DbConnectionString())));

services.AddBrighter()
.AddProducers(producers =>
{
producers.Outbox = new MySqlOutbox(outboxConfiguration);
producers.ConnectionProvider = typeof(MySqlConnectionProvider);
// Use the EF Core transaction provider with your DbContext
producers.TransactionProvider = typeof(MySqlEntityFrameworkTransactionProvider<GreetingsEntityGateway>);
})
.UseOutboxSweeper()
...
{
// ... your producer registry
producers.Outbox = new MySqlOutbox(outboxConfiguration);
producers.ConnectionProvider = typeof(MySqlConnectionProvider);
// Use the EF Core transaction provider with your DbContext
producers.TransactionProvider = typeof(MySqlEntityFrameworkTransactionProvider<GreetingsEntityGateway>);
})
.UseOutboxSweeper();
}

```

In our handler we take a dependency on our EF Core Context (derived from Db context). We explicitly start a transaction within the handler, because the Outbox is not within the Db Context we cannot rely on the DBContext's implicit transaction.
Expand Down
4 changes: 2 additions & 2 deletions contents/FeatureSwitches.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ By adding the **FeatureSwitch** Attribute or **FeatureSwitchAsync** Attribute, y
Registry**, [creating of which is described
later](/contents/FeatureSwitches.md#building-a-config-for-feature-switches-with-fluentconfigregistrybuilder).

In the following example, **MyFeatureSwitchedHandler** will only be run if it has been configured in the **Feature Switch Registry** and set to **FeatureSwitchStatus.On**.
In the following example, **MyFeatureSwitchedHandler** runs when the **Feature Switch Registry** sets it to **FeatureSwitchStatus.On**, and is skipped when it sets it to **FeatureSwitchStatus.Off**. With no registry at all, **FeatureSwitchStatus.Config** behaves as **On**. With a registry that has no entry for the handler, its **MissingConfigStrategy** decides; the **FluentConfigRegistryBuilder**'s default throws a `ConfigurationException`.

```csharp
using Paramore.Brighter;
Expand Down Expand Up @@ -64,7 +64,7 @@ public class MyIncompleteHandlerAsync : RequestHandlerAsync<MyCommand>

By default, when a feature switch is **Off**, the handler is skipped and the message is silently acknowledged and discarded. This is fine when you are using the Command Processor directly, but when consuming messages from an [External Bus](/contents/DispatchingARequest.md) you may want to hold messages on the channel until the feature is re-enabled, rather than losing them.

The `dontAck` parameter controls this behavior. When set to `true` and the feature is off, the attribute throws a `DontAckAction` instead of silently consuming the message. The [message pump](/contents/HowServiceActivatorWorks.md) leaves the message unacknowledged on the channel, and the transport re-delivers it after its visibility timeout expires.
The `dontAck` parameter controls this behavior. When set to `true` and the feature is off, the attribute throws a `DontAckAction` instead of silently consuming the message. The [message pump](/contents/HowServiceActivatorWorks.md) returns the message to the channel unacknowledged, and the transport delivers it again: the in-memory transport at once, a broker by its own redelivery rules.

```csharp
using Paramore.Brighter;
Expand Down
4 changes: 2 additions & 2 deletions contents/HandlerFailure.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ The right strategy depends on why processing failed and what you want to happen
- **Transient error, retry immediately:** Use `[UseResiliencePipeline]` with a Polly retry policy to retry within the same message pump cycle. If all retries fail, the exception propagates to the pump. See [Retry and Circuit Breaker](/contents/PolicyRetryAndCircuitBreaker.md).
- **Transient error, retry later:** Use `[DeferMessageOnError]` to catch any unhandled exception and requeue the message on the External Bus with a delay. The message becomes available to any consumer after the delay expires. For fine-grained control, throw `DeferMessageAction` directly. See [Requeue with Delay](#requeue-with-delay-defermessageaction).
- **Non-transient error, preserve for investigation:** Throw `RejectMessageAction` to route the message to a Dead Letter Queue (DLQ). Or use `RejectMessageOnErrorAttribute` as a backstop to catch any unhandled exception and reject. See [Reject to Dead Letter Queue](#reject-to-dead-letter-queue-rejectmessageaction).
- **Temporary block, try again after transport timeout:** Throw `DontAckAction` to leave the message on the channel. The transport re-delivers it after its visibility timeout expires. Or use `[DontAckOnError]` as a backstop. See [Don't Acknowledge](#dont-acknowledge-dontackaction).
- **Temporary block, try again after transport timeout:** Throw `DontAckAction` to leave the message on the channel. The pump returns it unacknowledged, and the transport delivers it again: the in-memory transport at once, a broker by its own redelivery rules. Or use `[DontAckOnError]` as a backstop. See [Don't Acknowledge](#dont-acknowledge-dontackaction).
- **Deserialization failure:** Throw `InvalidMessageAction` from a message mapper to route the message to an invalid message channel, keeping it separate from processing errors. See [Invalid Message Handling](#invalid-message-handling-invalidmessageaction).
- **Compensating action before failing:** Use `[FallbackPolicy]` to run cleanup or compensating logic when the handler fails, then let the exception propagate. See [Fallback Handlers](/contents/PolicyFallback.md).
- **Default (do nothing):** Let the exception propagate. The message is acknowledged and discarded. This is appropriate when errors are non-transient and you rely on logs and traces for investigation.
Expand Down Expand Up @@ -244,7 +244,7 @@ See [Backstop Attributes](#backstop-attributes) for pipeline ordering guidance a

### What Don't Acknowledge Does

Throwing `DontAckAction` tells the message pump to leave the message unacknowledged on the channel. The transport re-delivers it after its visibility timeout expires. A configurable delay (`DontAckDelay`, default 1 second) pauses the pump before processing the next message, preventing tight-loop CPU burn when a message is repeatedly not acknowledged.
Throwing `DontAckAction` tells the message pump to return the message to the channel unacknowledged (`Nack`), and the transport delivers it again: the in-memory transport at once, a broker by its own redelivery rules. A configurable delay (`DontAckDelay`, default 1 second) pauses the pump before processing the next message, preventing tight-loop CPU burn when a message is repeatedly not acknowledged.

Each `DontAckAction` increments the unacceptable message counter. If the counter reaches the configured `UnacceptableMessageLimit`, the pump shuts down. You can prevent shutdown by setting the `UnacceptableMessageLimit` to 0, or negative. You can also use `UnacceptableMessageLimitWindow` to control the period in which the limit is evaluated. This allows you to shut down for a burst of failures - typical if you have a poison pill message - but ignore failures that occur over time. (See below for more.)

Expand Down
9 changes: 7 additions & 2 deletions contents/HowConfiguringTheCommandProcessorWorks.md
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,11 @@ resiliencePipelineRegistry.TryAddBuilder(Globals.MYCIRCUITBREAKERANDRETRY,

When you attribute your code, you then use the key to attach a specific resilience pipeline:

``` csharp
```csharp
using Paramore.Brighter;
using Paramore.Brighter.Logging.Attributes;
using Paramore.Brighter.Policies.Attributes;

[RequestLogging(step: 1, timing: HandlerTiming.Before)]
[UseResiliencePipeline(Globals.MYRETRYPIPELINE, step: 2)]
public override TaskReminderCommand Handle(TaskReminderCommand command)
Expand All @@ -152,7 +156,8 @@ public override TaskReminderCommand Handle(TaskReminderCommand command)
A handler method takes only one `[UseResiliencePipeline]`; a second on the same method does not compile. If you need several strategies, such as a circuit breaker around a retry, compose them in one pipeline and attach that. The first strategy you add is the outermost; see [Combining Multiple Strategies](/contents/PolicyRetryAndCircuitBreaker.md#combining-multiple-strategies).

```csharp
// ...
using Paramore.Brighter.Policies.Attributes;

[UseResiliencePipeline(Globals.MYCIRCUITBREAKERANDRETRY, step: 1)]
public override TaskReminderCommand Handle(TaskReminderCommand command)
{
Expand Down
5 changes: 2 additions & 3 deletions contents/InMemoryOutbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,13 +87,11 @@ services.AddBrighter(options =>
```csharp
using System.Threading;
using System.Threading.Tasks;
using System.Transactions;
using Paramore.Brighter;

public class CreatePersonHandler : RequestHandlerAsync<CreatePerson>
{
private readonly IAmACommandProcessor _commandProcessor;
private readonly IAmAnOutboxAsync<Message, CommittableTransaction> _outbox;
private readonly PersonRepository _repository;

public override async Task<CreatePerson> HandleAsync(
Expand All @@ -104,7 +102,8 @@ public class CreatePersonHandler : RequestHandlerAsync<CreatePerson>
var person = new Person(command.Name, command.Email);
await _repository.SaveAsync(person);

// Deposit message to outbox (held in memory)
// Writes the message to the in-memory outbox and dispatches it at once;
// the entry stays, marked dispatched, until its time-to-live expires
await _commandProcessor.PostAsync(new PersonCreated { PersonId = person.Id }, cancellationToken: cancellationToken);

return await base.HandleAsync(command, cancellationToken);
Expand Down
2 changes: 1 addition & 1 deletion contents/KafkaConfiguration.md
Original file line number Diff line number Diff line change
Expand Up @@ -689,7 +689,7 @@ It is worth noting the following aspects of the code sample below:
* We provide two helpers, though you can pass your own settings if you prefer:
* **ConfluentJsonSerializationConfig.SerdesJsonSerializerConfig()** offers default settings for JSON serialization (many of these are passed through to Json.NET).
* **ConfluentJsonSerializationConfig.NJsonSchemaGeneratorSettings()** offers default settings for JSON Schema generation (such as using camelCase).
* The serializer writes a magic byte and the schema id ahead of the JSON, so the body is binary: we pass the bytes with an `application/octet-stream` **ContentType** and **CharacterEncoding.Raw**, as a round-trip through a UTF-8 string would corrupt the header. See [Message Mappers](/contents/MessageMappers.md).
* The serializer writes a magic byte and the schema id ahead of the JSON, so the body is binary: we pass the bytes with an `application/octet-stream` **ContentType**, and they reach Kafka as they are. A round-trip through a UTF-8 string would corrupt the header, and a relational Outbox with a text payload column makes one, so post through an Outbox configured with `binaryMessagePayload: true`. See [Message Mappers](/contents/MessageMappers.md).

``` csharp
using System.Net.Mime;
Expand Down
6 changes: 4 additions & 2 deletions contents/MessageMappers.md
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,9 @@ var payload = System.Text.Json.JsonSerializer.Serialize(request, new JsonSeriali
var body = new MessageBody(payload, new ContentType(MediaTypeNames.Application.Json), CharacterEncoding.UTF8);
```

If your payload is binary, then we provide two constructors that can be used to write bytes. For backwards compatibility these constructors also default to application/json and UTF-8. However, if you have binary content we recommend setting the media type to application/octet-stream and the character encoding to either **CharacterEncoding.Base64** if it needs transmission as a string, or **CharacterEncoding.Raw** if not).
If your payload is binary, then we provide two constructors that can be used to write bytes. For backwards compatibility these constructors also default to application/json and UTF-8; if you have binary content, set the media type to application/octet-stream. These constructors keep the bytes exactly as given, whatever **CharacterEncoding** you pass, and a transport sends those bytes: the encoding decides only how **Value** renders them as a string (below).

A string round-trip is what damages binary content, and a relational Outbox makes one when it stores the body in a text column: it writes **Value** and reads it back as UTF-8 text. Under **CharacterEncoding.UTF8** that corrupts any byte that is not valid UTF-8; under **CharacterEncoding.Raw** the consumer receives the base64 text instead of your bytes. Configure the Outbox with `binaryMessagePayload: true`, which stores the bytes themselves.

```csharp
public MessageBody(byte[]? bytes, ContentType? contentType = null, CharacterEncoding characterEncoding = CharacterEncoding.UTF8)
Expand All @@ -158,7 +160,7 @@ public MessageBody(in ReadOnlyMemory<byte> body, ContentType? contentType = null
// ...
```

For example, when writing a Kafka payload with leading bytes indicating the schema id, you would want to use a binary payload because conversion to and from a UTF8 string is lossy. Here we serialize the payload with the Kafka header (Magic Byte (0) + Schema Id Bytes) and a JSON payload using the Confluent Serdes serializer. Even though we serialize to JSON, because of the header bytes we treat the payload as binary:
For example, when writing a Kafka payload with leading bytes indicating the schema id, you would want to use a binary payload because conversion to and from a UTF8 string is lossy: a schema id byte of 0x80 or above does not survive it. Here we serialize the payload with the Kafka header (Magic Byte (0) + Schema Id Bytes) and a JSON payload using the Confluent Serdes serializer. Even though we serialize to JSON, because of the header bytes we treat the payload as binary:

```csharp
using System.Net.Mime;
Expand Down
Loading
Loading