Logging & Error Handling

Setting up SeriLog

To set up SeriLog to capture the desired logging information you will need to have the following NuGet packages installed:

To configure SeriLog, in the RegisterServices method in the service startup (Program.cs) add the following:

// Logging using Serilog
builder.Logging.AddSerilog();
var loggerOptions = new ConfigurationReaderOptions { SectionName = AuditConstants.AppSettingsSectionNames.Serilog };    
Log.Logger = new LoggerConfiguration()
                .ReadFrom.Configuration(builder.Configuration, loggerOptions)
                .Filter.ByExcluding("RequestPath like '/health%'")
                .Enrich.WithExceptionDetails()
                .Enrich.FromLogContext()
                .Enrich.WithSpan()                  
                .Enrich.With<ActivityEnricher>()
                .CreateLogger();

//Serilog.Debugging.SelfLog.Enable(Console.Error); 

You can add additional filters to exclude certain request paths you don’t want captured in logging (swagger for example).

Next add the following to appsettings.json

"Link:Audit:Logging:Serilog": {
  "Using": [ "Serilog.Sinks.Console", "Serilog.Sinks.Grafana.Loki" ],
  "MinimumLevel": {
    "Default": "Information",
    "Override": {
      "Microsoft": "Warning",
      "System": "Warning"
    }
  },
  "WriteTo": [
    { "Name": "Console" },
    {
      "Name": "GrafanaLoki",
      "Args": {
        "uri": "http://localhost:3100",
        "labels": [
          {
            "key": "app",
            "value": "Link-BoTW"
          },
          {
            "key": "component",
            "value": "Audit"
          }
        ],
        "propertiesAsLabels": [ "app", "component" ]
      }
    }
  ]
}

Additionally you can add the following to enhance the format of the console:

"Logging": {
    "LogLevel": {
      "Default": "Information",
      "Microsoft": "Warning",
      "System": "Warning"
    },
    "Console": {
      "FormatterName": "json",
      "FormatterOptions": {
        "SingleLine": true,
        "IncludeScopes": true,
        "TimestampFormat": "HH:mm:ss ",
        "UseUtcTimestamp": true,
        "JsonWriterOptions": {
          "Indented": true
        }
      }
    }
  }

Adding Problem Details

Problem details is an attempt at an industry standard to carry machine readable details of errors in a HTTP response to avoid the need to define new error response formats for HTTP APIs.

RFC-7807

This can be accomplished by using the .NET problem details service.

To enable to sending of Problem Details add the following in the RegisterServices method in the service startup (Program.cs):

//Add problem details
builder.Services.AddProblemDetails(options => {
    options.CustomizeProblemDetails = ctx =>
    {            
        ctx.ProblemDetails.Detail = "An error occured in our API. Please use the trace id when requesting assistence.";
        if (!ctx.ProblemDetails.Extensions.ContainsKey("traceId"))
        {
            string? traceId = Activity.Current?.Id ?? ctx.HttpContext.TraceIdentifier;
            ctx.ProblemDetails.Extensions.Add(new KeyValuePair<string, object?>("traceId", traceId));
        }      

        if (builder.Environment.IsDevelopment())
        {
            ctx.ProblemDetails.Extensions.Add("service", "Audit");
        }
        else 
        {
            ctx.ProblemDetails.Extensions.Remove("exception");
        }
        
    };                        
});

Add the following to the SetupMiddleware method in the service startup (Program.cs):

if (app.Environment.IsDevelopment() || app.Environment.EnvironmentName.Equals("Local", StringComparison.InvariantCultureIgnoreCase))
{
    app.UseDeveloperExceptionPage();
}
else
{
    app.UseExceptionHandler();
}

This will consume any 500 error returned by the API and turn it into a ProblemDetails model that is returned to the user. It will appear as follows:

{
    "type": "https://httpstatuses.io/500",
    "title": "Internal Server Error",
    "status": 500,
    "detail": "An error occured in our API. Please use the trace id when requesting assistence.",
    "traceId": "00-09f4ab47f959288745f3ef337ac773b7-fac718cf545d0a9d-01"
}

We can use the Trace Id to find relevant logs pertaining to the issue.

For this to work correctly you will need to return a 500 response type along with the exception in your API.

catch (Exception ex)
{
    _logger.LogError(new EventId(LoggingIds.GetItem, "Get Audit Event"), ex, "An exception occurred while attempting to retrieve an audit event with an id of {id}", id);
    throw;
}

Logging Ids

The AuditLoggingIds.GetItem in the above _logger.LogError method is a new option in .NET that allows you to assign a Logging Id to log types. In the example above it will give the log an id of 1002 from the following static class.

public static class AuditLoggingIds 
{
    public const int GenerateItems = 1000;
    public const int ListItems = 1001;
    public const int GetItem = 1002;
    public const int InsertItem = 1003;
    public const int UpdateItem = 1004;
    public const int DeleteItem = 1005;      
    public const int GetItemNotFound = 1006;
    public const int UpdateItemNotFound = 1007;
    public const int EventConsumer = 2000;
    public const int EventProducer = 2001;
    public const int HealthCheck = 9000;
}

This can allow for finding specific log types by the id and will show up in logs as follows:

UserScope Middleware

The LantanaGroup.Link.Shared.Application.Middleware middleware has been added to the shared project. This middleware will capture the username as well as the id of the user if an authenticated user is present in the .NET HttpContext.User.Identity object. This object will be populated based on if/how you have authentication enabled. For more information about enabled authentication see TBD.

This scoped user information will automatically be added to any logs created. To add the middleware simply add the following in the SetupMiddleware method in the service startup (Program.cs):

app.UseRouting();
app.UseCors("CorsPolicy");
app.UseAuthentication();
app.UseMiddleware<UserScopeMiddleware>();
app.UseAuthorization();    
app.UseEndpoints(endpoints => endpoints.MapControllers());   

This middlware should come after app.UseAuthentication(). When this is added the following will be included in a scope within your logs.

Adding data to exceptions

For additional context, you may add some additional data about the issue that caused the exception to be thrown as follows:

catch (NullReferenceException ex)
{
    _logger.LogDebug(new EventId(LoggingIds.GetItem, "Get Audit Event"), ex, "Failed to get audit event with an id of {id}.", id);
    var queryEx = new ApplicationException("Failed to get audit event", ex);
    queryEx.Data.Add("Id", id);
    throw queryEx;
}

The exception.Data.Add() method will add this information in the following locations:

Additional information about .NET logging can be found here.

LoggerMessage Attribute

Introduced in .NET 6, the LoggerMessage allows for performant logging.

public static partial class Logging
{
    [LoggerMessage(
      AuditLoggingIds.GenerateItems, 
      LogLevel.Information, 
      "New audit event created")]
    public static partial void LogAuditEventCreation(this ILogger logger, 
      [LogProperties]AuditEntity auditEvent);
}

The LogProperties attribute, requires .NET 8, will include all of the properties of the object in your log.

Kafka Consumer Error Handling — Dotnet Services

When consuming an event that cannot be deserialized, the consumer should catch the error through a ConsumeException. In the catch, the service should do the following:

  • Create an audit event to trace that an error occurred.
  • Produce an error event to the Error topic. Storing improperly structured events in an separate topic will allow for further investigate
  • Commit the consumer result to acknowledge back to the Kafka broker that the consumed event was processed.
flowchart LR
  SourceTopic["SOURCE TOPIC"] --> KafkaApp["KAFKA<br>APP"]
  KafkaApp -->|1| TargetTopic["TARGET TOPIC"]
  KafkaApp -->|2| ErrorTopic["ERROR TOPIC"]

  classDef source fill:#e6f4fa,stroke:#aaa;
  classDef target fill:#e6f4fa,stroke:#aaa;
  classDef error fill:#f8d7da,stroke:#c00,color:#000;

  class SourceTopic source;
  class TargetTopic target;
  class ErrorTopic error;

Example:

ConsumeResult<ReportScheduledKey, ReportScheduledValue> consumeResult;
try 
{
    consumeResult = _reportScheduledConsumer.Consume(cancellationToken);
}
catch (ConsumeException e)
{
    _logger.LogError($"Consume failure, potentially schema related: {e.Error.Reason}");
    var potentialFacilityId = Encoding.UTF8.GetString(e.ConsumerRecord.Message.Key);
    var potentialValue = Encoding.UTF8.GetString(e.ConsumerRecord.Message.Value);

    ProduceErrorEvent(potentialFacilityId, potentialValue);

    var auditValue = new AuditEventMessage
    {
        FacilityId = potentialFacilityId,
        Action = AuditEventType.Query,
        ServiceName = "QueryDispatch",
        EventDate = DateTime.UtcNow,
        Notes = $"Kafka ReportScheduled consume failure, potentially schema related \nException Message: {e.Error}",
    };

    ProduceAuditEvent(auditValue, e.ConsumerRecord.Message.Headers);

    _reportScheduledConsumer.Commit();

    continue;
}
private void ProduceErrorEvent(string key, string value) {
  var config = new ProducerConfig()
  {
      ClientId = "Error-QueryDispatch-ReportScheduled"
  };

  using (var producer = _errorProducerFactory.CreateProducer(config)) 
  {
      var headers = new Headers
      {
          new Header("X-Consumer", Encoding.UTF8.GetBytes("QueryDispatch")),
          new Header("X-Topic", Encoding.UTF8.GetBytes("ReportScheduledEvent"))
      };

      producer.Produce("Error", new Message<string, string> 
      { 
          Key = key,
          Value = value,
          Headers = headers
      });

      producer.Flush();
  }
}

Kafka Consumer Retry Handling

When a problem occurs while processing a consumed event (Facility not properly configured, API connection could not be established, could not connect to database, etc), we will want to attempt to perform retries to ensure that BotW processes each event successfully. When an exception is caught during the processing of an event, the service should do the following:

  • Create an audit event to trace that an error occurred.
  • Produce a retry event to the services Retry topic.
  • Commit the consumer result to acknowledge back to the Kafka broker that the consumed event was processed.
flowchart LR
  ReportScheduled --> QueryDispatch
  QueryDispatch --> Error
  QueryDispatch --> Retry
  QueryDispatch --> ConsumeSuccess
  Retry --> QueryDispatchRetry
  QueryDispatchRetry --> Retry
  QueryDispatchRetry --> ConsumeSuccess
  ConsumeSuccess:::final -->|Success or Retry| Retry

  classDef final fill:#d4af37,stroke:#333,color:white;

Kafka Consumer Error Handling — Java Services

How the Link Java services — MeasureEval and Validation — handle failures while consuming Kafka messages: what gets retried, what gets dead-lettered, on what schedule, and how to observe and test all of it. The mechanism is shared; each service plugs in its own topics.

This page is self-contained — it assumes familiarity with Kafka basics (topics, consumers, offsets) but nothing about the Link codebase.

Service Consumed topics (each with -Retry/-Error companions) Downstream event on success Per-message side resources
MeasureEval ResourcesNormalized, EvaluationRequested MeasureReportGenerated Resource cache (Redis/blob) — released only on success or terminal failure
Validation ReadyForValidation ValidationComplete None

The problem being solved

Link services communicate through Kafka: for example, when patient data has been normalized, a ResourcesNormalized message tells MeasureEval to evaluate the patient against a measure. Processing such a message can fail for two very different reasons:

  • Transient failures — a dependency is briefly down (MongoDB, Redis, blob storage, another service). The same message would succeed if simply tried again a bit later.
  • Permanent failures ("poison" messages) — the message content itself is bad (malformed JSON, an unparseable FHIR resource, a validation error). Retrying can never succeed.

Without deliberate handling, the first kind loses real patient data over a 30-second network blip, and the second kind can jam a consumer forever, re-failing on the same message. The design below retries transient failures automatically on a fixed schedule and moves permanent failures aside — into an error topic — where they are preserved for investigation instead of blocking everything behind them.

What happens when a message fails

Each consumed topic has two companion topics, created alongside it:

  • <Topic>-Retry — a holding area for messages awaiting another attempt. A dedicated listener re-reads them, waits until the message's scheduled due-time, and feeds it back through the exact same processing logic.
  • <Topic>-Error — the dead-letter topic. Messages that will never be retried again end up here, preserved with their original content and headers. Nothing consumes this topic; it is an archive for humans (and a future replay tool).
flowchart LR
  Main["ResourcesNormalized<br>(main topic)"] --> Consumer["Consumer<br>(processing logic)"]
  Consumer -->|success| Ack["done — offset committed"]
  Consumer -->|failure| Decide{"classify the<br>failure"}
  Decide -->|"transient +<br>attempts left"| Retry["ResourcesNormalized-Retry"]
  Decide -->|"poison, or attempts<br>exhausted, or retries disabled"| Error["ResourcesNormalized-Error<br>(dead letter)"]
  Retry -->|"after the backoff delay"| Consumer

  classDef error fill:#f8d7da,stroke:#c00,color:#000;
  class Error error;

Step by step, when processing throws an exception:

  1. The failure is classified (see the next section): can a retry ever help?
  2. Transient, attempts remaining → the message is republished to -Retry, stamped with its attempt count and an absolute "not before" timestamp (now + the configured delay). The retry listener redelivers it after that delay and the full processing logic runs again.
  3. Poison, attempts exhausted, or retries disabled → the message is republished to -Error. This is terminal: it will never be redelivered automatically.
  4. Only after the republish succeeds is the original message acknowledged. If the service crashes mid-flight — or the republish itself fails — the offset stays uncommitted and Kafka redelivers the original message. A message is never silently lost.

One special case: a message whose bytes can't even be deserialized fails before the processing logic runs at all. Such messages skip the retry ladder (deserialization can never succeed on retry) and go straight to -Error.

The diagram uses MeasureEval's ResourcesNormalized as the example; the flow is identical for its EvaluationRequested topic and for Validation's ReadyForValidation (ReadyForValidation-Retry / ReadyForValidation-Error).

One operational prerequisite: the -Retry and -Error topics are not auto-created — they must exist in the broker (they are part of the standard topic provisioning). In a fresh environment where they are missing, failed messages have nowhere to go.

Transient vs poison: how failures are classified

Outcome When Where the message goes
Retry Any exception not classified as poison, with attempts remaining -Retry, redelivered after the scheduled delay
Dead-letter: poison The exception (anywhere in its cause chain) is a content problem that can never succeed — in both services: FHIR parse errors, validation errors, message-format and deserialization errors -Error, immediately, without burning retry attempts
Dead-letter: exhausted All attempts used up -Error
Dead-letter: disabled Retries are turned off by configuration -Error on the first failure

The default is deliberately retry: anything not on the service's explicit poison list — including infrastructure outages like "cannot connect to Redis" — is treated as transient. The flip side of that rule for developers: processing code must let such exceptions escape unhandled. Catching and swallowing a Redis outage would not prevent the failure — it would convert it into a silently wrong result (e.g. evaluating an empty patient bundle and reporting a false "not reportable") with no retry and no dead letter.

The retry schedule

The delays come from configuration (next section). Both services deploy the same schedule:

Attempt Delay before it Elapsed since first failure
1 (original) — 0
2 20 s ~20 s
3 60 s ~1 min 20 s
4 (last) 120 s ~3 min 20 s
dead letter — ~3 min 20 s

So a dependency outage shorter than ~3.5 minutes self-heals with no data loss and no manual action: whichever attempt lands after the dependency returns completes the real work, produces the normal downstream event, and the message never reaches -Error.

Configuration

Two properties control everything; both live under spring.kafka.retry (in application.yml, overridable per environment via Azure App Configuration):

spring:
  kafka:
    retry:
      consumer-retry-duration: [PT20S, PT60S, PT120S]  # ISO-8601 durations
      disable-retry-consumer: false
  • consumer-retry-duration is the single source of truth for both how many retries happen and when: N entries mean N retries (N+1 total attempts), and entry K is the delay before retry K. An empty list means the first failure dead-letters. This mirrors the .NET services' ConsumerSettings:ConsumerRetryDuration, so both stacks are tuned the same way.
  • disable-retry-consumer: true sends every failure straight to -Error and doesn't start the retry listener at all.
  • Changing depth, timing, or turning retries off is a configuration change and a restart — no code change.

Guarantees — and their consequences

  • At-least-once delivery. A message is acknowledged only after it has been fully processed or durably parked in -Retry/-Error. The trade-off: a crash or restart at the wrong moment can deliver a message twice. This is by design (the alternative is losing it), and it means duplicate downstream events (e.g. two MeasureReportGenerated for one input) are possible and expected under failure conditions — downstream consumers must tolerate them.
  • Dead letters are terminal today. Nothing consumes -Error; there is no automatic replay. Recovering a dead-lettered message is currently a manual operation. (A cross-service replay capability is a known future work item.)
  • Cached data survives retries. Messages like ResourcesNormalized don't carry the patient data themselves — it sits in a shared resource cache (Redis / blob storage) keyed by correlation id. The rule: a failed attempt must leave the cache untouched, because the retry needs it. The cache is released at exactly two moments — successful processing, or after the message is durably dead-lettered — never in between. Validation carries no per-message side resources, so it has no cleanup step at all.
  • Dependency clients fail fast on the consumer path. Blob-storage clients used during message processing are capped at 2 quick tries (~10 s) instead of the SDK default of ~3 minutes of hidden internal retrying. One system owns retries — the visible, configurable Kafka ladder — so outages surface in seconds and the schedule above is the real schedule.

For QA: what to expect and how to verify

Expected behavior by scenario

Scenario Expected behavior
Malformed / poison message Goes to -Error immediately (no retry attempts); normal traffic continues unaffected
Dependency down briefly (< ~3.5 min) Retry ladder runs (20 s / 60 s / 120 s); once the dependency is restored, the next attempt succeeds and produces the normal downstream event; -Error stays empty
Dependency down past the ladder (~3.5 min+) Message lands in -Error after 4 attempts; downstream event is not produced
Service restarted mid-ladder Message is redelivered and continues; a duplicate downstream event is possible and acceptable
Retries disabled via config First failure goes straight to -Error

Which dependency outages actually trigger the ladder differs per service — the ladder only engages when a failure makes processing throw, and each service has different hard dependencies:

  • MeasureEval retries on: Redis down (its resource cache is a hard dependency), MongoDB down, blob storage (Azurite locally) down. Data-only triggers, needing no infrastructure change: an EvaluationRequested referencing a nonexistent previous report, or an unknown measure id.
  • Validation retries on: blob storage down while the message carries a valid payloadUri, or its SQL Server database down (Validation has no MongoDB dependency).
  • Validation trap: stopping Redis does not produce retries in Validation. Its Redis layer is a read-through cache whose errors are deliberately swallowed and treated as cache misses — processing continues against the backing store. Expect warning logs and a degraded health check, not -Retry traffic. This is correct behavior, not a bug.

What to observe

Three signals tell the whole story per test message:

  1. Service logs — each retry hop and each dead-letter decision is logged. Against local docker-compose (link-measureeval / link-validation):
    docker logs link-measureeval -f | grep -E "Retry attempt|Max retry|Cache cleanup"
    docker logs link-validation  -f | grep -E "Retry attempt|Max retry"
    
  2. The -Error topic offset — did the dead-letter count go up? (Substitute ReadyForValidation-Error when testing Validation.)
    docker exec kafka-broker /opt/bitnami/kafka/bin/kafka-get-offsets.sh \
      --bootstrap-server kafka_b:9092 --command-config /tmp/client.properties \
      --topic ResourcesNormalized-Error
    
  3. The downstream event — MeasureReportGenerated (MeasureEval) / ValidationComplete (Validation) must appear on success and only on success.

Inducing failures: stopping dependency containers

With the docker-compose stack running, use the levers from the list above: docker stop mongo (MeasureEval), docker stop mssql (Validation), docker stop azurite (MeasureEval's resource cache; Validation's payload reads), docker stop redis_cache (MeasureEval only — see the Validation trap above). Restart the dependency mid-ladder to verify recovery.

docker stop and docker pause are two different scenarios, and both are worth running: stop gives connection-refused (fast, clean failures); pause gives a TCP black hole — connections hang until a timeout fires — which exercises the configured timeout bounds (the 2 s Redis command timeout, Mongo's serverSelectionTimeoutMS) and is closer to a real network outage.

Inducing failures: breaking configuration

Config-based faults complement container stops: they model real misconfiguration incidents, and they target one service without disturbing other consumers of the same dependency.

Mechanics (local docker-compose): add or override the property in the service's environment: block in docker-compose.yml, in dotted form exactly like the existing entries (e.g. spring.data.mongodb.uri: ...), then recreate the service with docker compose up -d measureeval (or validation). In deployed environments the same property paths go into Azure App Configuration plus a service restart.

The golden rule: a broken config only exercises the retry ladder if the service still boots and the failure happens at message-processing time. Some breakages fail at startup instead — a crashed or never-healthy container runs no consumers and produces zero retries, which reads like a failed test when it's really a failed setup. The traps are flagged below.

MeasureEval:

Lever Property to break Behavior
Redis port spring.data.redis.port: 6380 (valid but unused port) ✅ Boots fine; every cache read throws at processing time → ladder. Use a valid-but-dead port — an out-of-range value like 66666 fails property binding and kills startup
Redis password spring.data.redis.password: wrong ✅ Runtime (the client connects lazily) — auth failure at first cache read → ladder
Mongo credentials spring.data.mongodb.uri: mongodb://linkuser:wrongpass@mongo:27017/link-measureeval?authSource=admin&serverSelectionTimeoutMS=4000 ✅ Runtime — every Mongo access fails with AuthenticationFailed (error 18) → ladder. Works even though the local compose Mongo has no auth enabled: the server still validates the SCRAM handshake against its user catalog, and the user doesn't exist
Mongo unreachable dead port in the URI, e.g. mongodb://mongo:27018/...?serverSelectionTimeoutMS=4000 ✅ Runtime; equivalent to stopping Mongo but scoped to this one service
Blob container resource-cache.blob-storage.blob-container-name: no-such-container ✅ Runtime — cache blob reads 404 (ContainerNotFound) → ladder; the fail-fast client (2 tries, 10 s cap) surfaces it in seconds. Breaking the account key inside the connection string gives the 403 variant

Mongo URI rules that make or break the test: keep the URI well-formed (a malformed one — e.g. a port inside an +srv URI — is a startup crash, not a runtime failure), and always keep serverSelectionTimeoutMS=4000 — the driver's default is 30 s of server selection per attempt, which stalls every rung and makes the ladder look hung.

Validation:

Lever Property to break Behavior
SQL Server password spring.datasource.password: wrong ❌ Startup killer, not a lever — Hibernate opens a connection during boot (ddl-auto: update), so the container never becomes healthy and no consumer runs. For DB-driven retries, break the DB after startup: docker stop mssql or docker pause mssql
Blob container internal-blob-storage.blob-container-name: no-such-container ✅ Runtime — the report-payload read fails (for messages with a valid payloadUri) → ladder
Redis (any breakage) port, password, or stopped ❌ Deliberately no retries — see the Validation trap above (cache errors are swallowed as misses)

Config faults are sticky — know which command heals them. A container keeps the environment it was created with: reverting docker-compose.yml does not change the running service, so the fault stays live until the container is recreated. The two commands do opposite things:

Command Effect
docker restart link-measureeval Keeps the injected fault (restart preserves the container's env)
docker compose up -d measureeval Removes it (recreates the container from the current file) — this is also the "dependency restored mid-ladder" recovery step; the in-flight message is parked in -Retry and survives the recreation

Two observation caveats while a config fault is live: the container may show unhealthy purely from its health-check probes hitting the broken dependency — that is not ladder traffic, and the resulting log noise (on HTTP threads) is not retry activity. Ladder behavior only appears once a message is actually consumed, so always produce a test record and watch for the Retry attempt log lines.

Producing test records

A poison message is any record whose value isn't valid JSON for the topic's type. To produce test records, use the Kafka console producer inside the broker container from Git Bash, not PowerShell (PowerShell pipes prepend a byte-order mark that corrupts the message key):

printf '%s\n' '{"facilityId":"MyFacility","patientId":"p1"}|{"queryType":"INITIAL",...}' |
  MSYS_NO_PATHCONV=1 docker exec -i kafka-broker \
  /opt/bitnami/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka_b:9092 --producer.config /tmp/client.properties \
  --topic ResourcesNormalized --property parse.key=true --property 'key.separator=|'

For developers: how it's built and how to adopt it

The mechanism lives in the shared Java module and is service-agnostic; each service plugs in its own topics and its own poison list.

Background on why it's hand-rolled: these services process messages on a background executor (long-running CQL evaluation would otherwise exceed Kafka's poll interval and trigger consumer rebalancing). Exceptions therefore surface on the executor thread, which Spring Kafka's native retry machinery cannot see — so the services own the retry policy (classification, attempt counting, schedule, destination) and reuse Spring only for the retry plumbing (the listener that delays -Retry redeliveries).

Shared component Role
AbstractAsyncConsumer Executor hand-off; ack on success, route-then-ack on failure
RetryTopicRecoverer Makes the routing decision, stamps attempt/due-time headers, publishes
RetryTopicRecovererFactory Assembles the above; takes the service's poison set and an optional terminal-cleanup hook as parameters
KafkaRetryConfig Binds the two spring.kafka.retry properties; derives attempt count and per-attempt delay

Relationships

flowchart LR
nlogging_error_handling_6DA934E6["Design: Logging & Error Handling"]