Skip to content

Concurrency Pipeline

Use the Concurrency Pipeline to automatically retry a Command when its execution returns a concurrency conflict.

The pipeline is useful with optimistic concurrency, where two operations can attempt to update the same resource and one operation may fail because the resource version has changed.

The concurrency retry configuration belongs to the command bus. It does not apply to queries.

You need:

  • an AppBuilder instance;
  • a Command and its handler;
  • concurrency conflicts represented as a failed Result with a conflict error;
  • the command executed through the Xeno.JS mediator.

If you are creating the Command itself, see Create a Command.

Configure commandBus.concurrency inside addPipeline():

const app = new AppBuilder()
.addPipeline((config) => {
config.commandBus.concurrency = {
maxRetries: 3,
delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
},
}
})
.build()

Once configured, Xeno.JS adds the concurrency retry behavior to the Command pipeline.

You do not need to register ConcurrencyRetryPipeline manually.

Use maxRetries to control the maximum number of times the command pipeline is executed.

For example:

config.commandBus.concurrency = {
maxRetries: 3,
delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
},
}

With maxRetries: 3, a command can be executed up to three times when it keeps returning a concurrency conflict:

Attempt 1
│
├── success ───────────────► return result
│
└── conflict
│
▼
delay
│
▼
Attempt 2
│
├── success ───────────────► return result
│
└── conflict
│
▼
delay
│
▼
Attempt 3
│
├── success ───────────────► return result
│
└── conflict ───────────────► return failed Result

maxRetries must be a positive integer.

A value of 0 or a negative value is invalid.

Use delayConfig to configure the delay applied between concurrency-conflict attempts:

config.commandBus.concurrency = {
maxRetries: 5,
delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
},
}

baseDelayMs defines the base delay in milliseconds before another attempt.

It must be a non-negative integer.

baseDelayMs: 100

maxJitterMs defines the maximum jitter added to the retry delay.

It must be a non-negative integer.

maxJitterMs: 50

Using jitter helps avoid multiple concurrent operations retrying at exactly the same time.

The concurrency configuration is optional.

You can enable it with the defaults:

const app = new AppBuilder()
.addPipeline((config) => {
config.commandBus.concurrency = {}
})
.build()

When a value is omitted, Xeno.JS uses the pipeline defaults.

You can override only the retry count:

const app = new AppBuilder()
.addPipeline((config) => {
config.commandBus.concurrency = {
maxRetries: 5,
}
})
.build()

Or provide the complete configuration:

const app = new AppBuilder()
.addPipeline((config) => {
config.commandBus.concurrency = {
maxRetries: 5,
delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
},
}
})
.build()

The retry pipeline does not retry every failed command.

A result is retried only when its error is a conflict with:

  • the CONFLICT error code;
  • HTTP status 409.

For example, a command handler can return a conflict result when an optimistic concurrency check detects that the entity was modified by another operation.

return Result.fail(
AppError.conflict(
request.intent,
'The user was modified by another request.',
),
)

That failed result is eligible for the concurrency retry pipeline.

Other failures are returned immediately.

For example:

Command
│
▼
Handler
│
├── success ───────────────► return result
│
├── validation error ──────► return error
│
├── authorization error ───► return error
│
└── conflict (409) ────────► retry

When a command returns a concurrency conflict:

  1. Xeno.JS checks whether the error is a conflict;
  2. if the retry limit has not been reached, it waits using the configured delay and jitter;
  3. the command pipeline is executed again;
  4. the new result is evaluated again.

A successful retry immediately returns the successful result.

For example:

Command
│
▼
Conflict
│
▼
wait
│
▼
Command again
│
▼
Success
│
▼
Return success

A non-concurrency error is not retried.

If every allowed attempt returns a concurrency conflict, the pipeline stops retrying and returns a failed conflict Result.

The returned error indicates that the maximum retry attempts were exceeded.

For example:

Maximum retry attempts (3) exceeded due to concurrency conflicts.

The command is not executed again after the limit has been reached.

A typical application can configure the command bus as follows:

import { AppBuilder } from '@xeno-js/shared'
const app = new AppBuilder()
.addPipeline((config) => {
config.commandBus.concurrency = {
maxRetries: 3,
delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
},
}
})
.build()

Your command handler is responsible for returning a conflict when an optimistic concurrency check fails:

import { AppError, Result } from '@xeno-js/shared'
export class UpdateUserHandler extends BaseHandler<
UpdateUserCommand,
UpdateUserResult
> {
protected async executeAsync(
request: UpdateUserCommand,
): Promise<Result<UpdateUserResult>> {
const user = await this.userRepository.findById(request.userId)
if (!user) {
return Result.fail(
AppError.notFound(
request.intent,
'User not found.',
),
)
}
if (user.version !== request.version) {
return Result.fail(
AppError.conflict(
request.intent,
'The user was modified by another request.',
),
)
}
// Update the entity and persist it.
return Result.ok({
userId: user.id,
})
}
}

With concurrency retries enabled, a conflict can cause the command to be executed again according to the configured retry policy.

The retry configuration is part of:

config.commandBus.concurrency

It is therefore applied to Commands.

It is not part of:

config.queryBus

and does not automatically retry Queries.

The application pipeline is conceptually:

Command
│
▼
Common pipelines
│
▼
Command-specific pipelines
│
└── Concurrency Retry
│
▼
Handler

If you need query result caching or other query-specific behavior, see the corresponding query pipeline documentation.

Validation and authorization errors are not retried

Section titled “Validation and authorization errors are not retried”

Concurrency retry is intentionally limited to conflict errors.

For example:

Validation failure
│
└──► return immediately
Authorization failure
│
└──► return immediately
Not found
│
└──► return immediately
Conflict / 409
│
└──► retry

This means that increasing maxRetries does not cause unrelated application errors to be retried.

Check that concurrency is configured under commandBus:

config.commandBus.concurrency = {
maxRetries: 3,
}

Then check that the handler returns a conflict error.

The retry pipeline only recognizes errors with:

  • conflict error code;
  • HTTP status 409.

A generic failed Result will not trigger a retry.

Check the error returned by the handler.

For example, this is a concurrency conflict:

return Result.fail(
AppError.conflict(
request.intent,
'Concurrent update detected.',
),
)

Whereas returning another error type will not activate the retry behavior.

maxRetries must be a positive integer.

Invalid:

maxRetries: 0

Invalid:

maxRetries: -1

Valid:

maxRetries: 3

Both values must be non-negative integers.

Invalid:

delayConfig: {
baseDelayMs: -100,
maxJitterMs: 50,
}

Valid:

delayConfig: {
baseDelayMs: 100,
maxJitterMs: 50,
}

Xeno.JS is an MIT-licensed open source project. It can grow thanks to the support of these awesome people. If you’d like to join them, please read more at support section