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
490 changes: 33 additions & 457 deletions README.md

Large diffs are not rendered by default.

186 changes: 72 additions & 114 deletions RedShirt.Example.JobWorker.slnx

Large diffs are not rendered by default.

51 changes: 51 additions & 0 deletions docs/health-probes.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
# Health probes

Notes on health probes.

## Overview

The JobWorker application has a configurable set of health pages. When health endpoints are enabled, health is currently
determined by the amount of time since the most recent major exception caught in the `RedShirt.Example.JobWorker.Core`
project when it interacts with a service from the `RedShirt.Example.JobWorker.Common.Distributed` project or the
`IJobSource` implementation in one of the `JobManagement` projects.

| Endpoint | Purpose | Healthy response | Unhealthy response |
|-------------------|------------|------------------------|------------------------------|
| `GET /live` | Liveness | `200` plain text `OK` | N/A |
| `GET /health` | Health | `200` plain text `OK` | `503` plain text `unhealthy` |
| `GET /statistics` | Statistics | `200` JSON (see below) | N/A |

Environment variables related to health:

* `HEALTH__ENABLED`: HTTP listener with health pages (default: `true`). When `false`, the worker runs without binding a
health port.
* `HEALTH__PORT`: TCP port for health endpoints, bound on `0.0.0.0` (default: `8080`).
* `HEALTH__RECENT_INCIDENT_THRESHOLD_SECONDS`: Amount of seconds after a major exception in `Core` project for which the
system will be considered unhealthy.
* `JOBS__HALT_ON_FAILURE`: Related. If set to `true`, then the application shall immediately throw major exceptions to
crash the application, making the health system moot. Only recommended for local development.

## Statistics Example

This is an example of the returned statistics model (C# definitions can be found in
`RedShirt.Example.JobWorker.Common.Health` in `Models/StatisticsModel.cs`:

```json
{
"lifetime": {
"successfulTimings": {
"average": "00:00:00",
"max": "00:00:00",
"min": "00:00:00"
},
"totals": {
"received": 0,
"successful": 0,
"cancelled": 0,
"failed": 0,
"invalidData": 0
}
},
"uptime": "00:12:34.5678900"
}
```
193 changes: 193 additions & 0 deletions docs/idempotency.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,193 @@
# Idempotency

## Overview

In order to properly implement the idempotent consumer pattern, the outcome of processing the same message repeatedly
must be the same as processing the message once.

This template has support for idempotent operations by way of Redis caches. Idempotency is based off of the message ID
set in the messaging source, which is internally referred to as the idempotency ID. The idempotency ID is also exposed
to the `IJobLogicRunner` implementation so that implementation logic can potentially make decisions based on this
identifier (this is encouraged if possible, see the [Idempotency in Job Logic](#idempotency-in-job-logic) section of
this document).

## Feature History

Prior to its implementation in this template, the idea for idempotency came from working on an unrelated RabbitMQ
message consumer application. The connection to RabbitMQ would become unstable when the server was under load. Because
of how RabbitMQ functions, this meant that ownership of the messages would be lost by the consumer. Even in an
environment with only one message consumer instance, the retrieval of a message is tied to the specific connection that
it came from. Even if the connection to RabbitMQ is reestablished, it cannot be used to acknowledge the message from the
previous connection.

Thus, the idempotency system in this template and the other worker had two main objectives:

* The same message from a source should not be processed twice.
* The same message from a source should *especially* not be processed twice simultaneously.
* If a message is received twice, then that suggests that the previous retrieval of the message lost custody and is
unable to acknowledge. The retrieval with custody should use the response status of the retrieval that lost custody
once it becomes available rather than re-process the message.

## Feature Architecture

* Idempotency operations are centralized in the Idempotency Execution Service (`IIdempotencyExecutionService`).
* The Idempotency Execution Service is used extensively by the executor worker threads. However, the idempotency service
is also written to return permissive stand-ins in the event that idempotency is not enabled in configuration.
* If idempotency support is enabled, then a monitor worker thread shall occasionally check the status of messages that
were detected as being re-received in parallel to an original receive.
* This parallel running could happen on this instance of the job worker process or another instance, as judgment is
made by a distributed lock based in a Redis instance.
* The monitor thread shall periodically try to re-acquire an exclusive lock on the message ID and to use the cached
result to immediately acknowledge the message based on past processing.
* If the monitor thread is able to acquire a lock but fails to retrieve a cached result, then the message is
re-flagged in the job repository as being a candidate for execution.

### Redis Connection Instability Tolerance

Though being a proper idempotent consumer is the overall goal, this general template prioritizes overall stability over
strict idempotency. Non-critical exceptions encountered while interacting with Redis at the low-level are captured by a
safety layer.

If the application fails to interact with Redis, then the safety layers will enter a "disgrace" state, in which the
lower-level Redis services will not be attempted until the "disgrace" period has passed.

## Idempotency IDs

The Idempotency ID of a message is its unique identifier that allows the idempotency system to function. The application
will not crash if receives a job with a null idempotency ID, but it won't be able to act as an idempotent consumer.

* For message brokers like SQS or Azure Service Bus, the Idempotency ID value is set off of the messages ID from the
system.
* For more stream-like job sources such as Kinesis or Kafka, the Idempotency ID value is based on an indication of a
record's position in the stream.

For many job sources and configurations, this identifier is automatically generated. However, there are some sources and
configurations where it is not set.

### Concerning Idempotency ID Uniqueness

Many message sources can automatically provide Idempotency IDs that are reliably unique. However, some services allow
them to be specified by the publisher submitting the message.

The Idempotency IDs are considered to be reliably unique based on of the configuration variable
`JOBS__IDEMPOTENCY__IDEMPOTENCY_IDS_CAN_REPEAT=false`. If the IDs are said to not repeat, then a successful
acknowledgement of a message shall mean that the cached result for that message will be cleared or not entered into the
cache at all. This is done in hope of saving cache resources.

### RabbitMQ Message IDs

Of the current roster of job sources, RabbitMQ has no option to automatically generate a message ID for the application
to take as an idempotency key. If you are using RabbitMQ and wish to make use of idempotency, then you will need to make
sure that your message publishers are providing a message ID.

In the RabbitMQ browser view, this can be done by manually specifying the `message_id` property.

In C#, this would look like this:

```csharp
var properties = channel.CreateBasicProperties();
properties.MessageId = Guid.NewGuid().ToString();
properties.Persistent = true; // Optional: make message persistent

var body = Encoding.UTF8.GetBytes("Hello RabbitMQ");

// Assume that channel has already been declared.
channel.BasicPublish(
exchange: "my-exchange",
routingKey: "my-routing-key",
mandatory: false,
basicProperties: properties,
body: body);
```

In Python (using the `pika` module):

```python
#!/usr/bin/env python

import uuid
import pika

properties = pika.BasicProperties(
message_id=str(uuid.uuid4()),
content_type="application/json",
delivery_mode=2, # Optional: make message persistent
)

body = b"Hello RabbitMQ"

# Assume that channel has already been declared.
channel.basic_publish(
exchange="my-exchange",
routing_key="my-routing-key",
body=body,
properties=properties,
)
```

### Redis Streams Message IDs

In practice, Redis Streams seems to also require manual setting of a message ID.

Important distinction for this template: the Redis stream entry ID is exposed as the job's `MessageId`, but the
idempotency key is taken from a `message_id` **field** on the stream entry. If you are using Redis Streams and wish to
make use of this template's idempotency features, then your publishers should set that field. Auto-generating or
manually specifying the Redis stream entry ID alone is not enough for the idempotency system.

Example in C# (StackExchange.Redis):

```csharp
var db = multiplexer.GetDatabase();
var fields = new NameValueEntry[]
{
new("body", """{"SleepDurationSeconds":12}"""),
// Supply a specific Redis stream entry ID
new("message_id", Guid.NewGuid().ToString()) // Idempotency ID for this template
};

var specificEntryId = await db.StreamAddAsync("jobs", fields);
```

In Python (`redis-py`):

```python
#!/usr/bin/env python

import json
import uuid

import redis

client = redis.Redis(host="localhost", port=6379, decode_responses=True)
values = {
"body": json.dumps({"SleepDurationSeconds": 12}),
# Supply a specific Redis stream entry ID
"message_id": str(uuid.uuid4()), # Idempotency ID for this template
}

specific_entry_id = client.xadd("jobs", values)
```

Documentation purports that one can provide an asterisk to request that Redis auto-generate an ID for a message, but
this has not been my experience in practice.

## Idempotency in Job Logic

Being a general template, this application does not make use of the exposed idempotency ID in job logic, nor does it use
any other job properties in a similar manner. That being said, I would like to make a case for considering doing so.

Consider a message that requires a complex operation with multiple steps. For the sake of this document, we'll imagine
that there are two steps:

1. Submit information in remote system A.
2. Submit information in remote system B.

If a previous handling of the job previously failed on step 2, then there is a risk of repeating step 1. In our
perfect-world (or next to perfect, we are dealing with a re-run after all), we would not repeat step 1 by giving our
system the awareness that the step had been executed.

Possible options for being state-aware:

* Information could be placed in a cache such as Redis.
* To re-use our imaginary complex operation from earlier in this section, system A could potentially be polled to
confirm if the information had already been submitted.
72 changes: 72 additions & 0 deletions docs/job-source-notes-kafka.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
# Notes on Implementing Kafka

Notes on Kafka.

## Overview

Kafka is more of an append-only event log rather than a traditional message queue.

I go into more detail on each of these points below, but the cliff notes for implementing Kafka are:

* This template's version of Kafka assumes no authentication, which is currently left as an excercise for the reader.
* It is strongly advised to use Batch mode polling when using Kafka as a job source.
* It is strongly advised to enable idempotency handling for Kafka. Basic enabling of idempotency handling hinges off of
the `JOBS__IDEMPOTENCY__ENABLED` environment variable, with other options described in the configuration section of
this document and demonstrated in `test/local/docker-compose.yaml`.

## Kafka Authentication

This general template was tested against a local Kafka container with no authentication set up. In addition to that,
there are currently 5 different SASL mechanisms to choose from when implementing authentication. Implementing
authentication for Kafka when adapting this template is currently left as an exercise for the reader.

Kafka clients are constructed in the `RedShirt.Example.JobWorker.JobManagement.Kafka` project, in
`Factories/KafkaConsumerFactory.cs`.

## Kinesis Comparisons and Batch Mode Recommendation

A Kafka topic is very similar to a Kinesis stream. I am going to be comparing Kafka to Kinesis very heavily in this
section because Kinesis is a much more established job source implementation that I have more experience with. For notes
specifically on Kinesis, please refer to [`job-source-notes-kinesis.md`](job-source-notes-kinesis.md).

The Kinesis comparison carries down to a basic implementation level of the technology. A Kafka topic is divided into
partitions, just as a Kinesis stream is divided into shards. Processing jobs from either of these sources involves some
layer of the process managing shard/partition ownership

However, a major difference between the available interfaces for Kafka and Kinesis and their implementations in this
template is how ownership of a partition/shard works:

* In Kinesis, the job source's application code lists and iterates through shards in an attempt to find one that does
not have a distributed lock. The Kinesis job source then performs a `GetRecords` operation on that shard.
* In Kafka, our options are more limited.
* Kafka does have an option to list individual partitions, but this is considered more of an admin action.
* Instead, the client declares a Kafka consumer which simply calls `consumerObject.Consume(TimeSpan)`, with a
TimeSpan for timeouts.
* Along the same lines: with no ability to iterate through partitions, ownership of a partition is out of the
client's hands.
* The Kakfa server/cluster calculates ownership of a partition within a consumer group when the number of
partitions changes or the number of connected clients changes. This means that a Kafka client can lose access
to a partition while still working exactly as intended.

Commiting a message in Kafka (done by its offset) implies that every message before it in the partition has also been
processed. This template sorts the jobs retrieved by a job source and messages could be run in parallel worker threads
with different finish times. Under these conditions, it cannot be guaranteed that the batch of messages being commited
during acknowledgement is the next one on the partition's to-do list.

Because of this, it is strongly recommended to run message polling in Batch mode as opposed to Loader mode. With Batch
mode, the job source is polled as soon as the previous batch has finished. In Loader mode, the job loader handler could
have to wait for several seconds (based on the value of the `JOBS__MAX_IDLE_WAIT_SECONDS` environment variable) before
polling again.

## Kafka Idempotency Handling

As described in the above section on Kinesis comparisons, Kafka clients in a consumer group do not control over what
partitions they have authority to commit to. Ownership of a partition is reconsidered when the number of clients in a
consumer group or the number of partitions in a topic changes. A client can lose commit rights to a topic partition
through no fault of its own.

Because of the above point, it is *strongly* encouraged to enable idempotency support for your application if you are
using Kafka with multiple consumers. Basic enabling of idempotency handling hinges off of the
`JOBS__IDEMPOTENCY__ENABLED` environment variable, with other options described in the configuration section of this
document and demonstrated in `test/local/docker-compose.yaml`. For more information on the idempotency system, please
refer to [`idempotency.md`](idempotency.md).
41 changes: 41 additions & 0 deletions docs/job-source-notes-kinesis.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
# Notes on Implementing Kinesis

Notes on Kinesis.

## Overview

Kinesis is a bit different from most of the other job sources in this template.

A Kinesis stream is fundamentally composed of multiple shards, and the shards contain the job record messages. A default
Kinesis stream will have 4 shards.

Fundamentally, the Kinesis as a source of records is different from the other messaging technologies covered within this
template in a number of ways:

* Kinesis is a stream rather than a true message broker.
* As a stream, individual messages have no built-in mechanism to be reclaimed by the queue in the event of the container
being stopped by a sudden and catastrophic problem outside its control (e.g. hardware failure). This is why the job
worker will only proceed to move the tracker for a shard past the current batch of messages after all messages have
been attempted. It is an intentional safety mechanism.
* The stream stores messages individually, but the iterator string used to progress in a short-term context only
operates in batches.
* The stream stores messages in sequential order, but this template supports prioritizing messages received in an
arbitrary order as defined by the chosen implementation of `ISourceMessageSorter`.

## Kinesis AI Audit Notes

The solutions to the above considerations for Kinesis are part of why an AI audit of the key Kinesis using the Composer
model took (and sometimes continues to take) issue with many points of the Kinesis job source's design. These objections
include but are not limited to:

* Worrying about "leaking" locks by keeping them for later storage rather than encapsulating all remaining operations in
the method in a try-finally statement.
* Assuming that messages will not be acknowledged if the job failed to process, or that some message Ids will never be
acknowledged for some other reason.
* This one in particular might be solvable with a different method name that better implies that it is always
called, but this is not a priority.
* Worrying about the "all-or-nothing" nature of processing a batch of messages from a shard. Composer is under the
impression that progress can be incremented per message. It is wrong.

While Composer really doesn't like the Kinesis job source in particular, these issues are in fact fundamental to the
operation of Kinesis within this framework.
Loading