From 7bfa5744a287d2cb7b7bf748aad4b23306803192 Mon Sep 17 00:00:00 2001 From: intech Date: Sat, 4 Jul 2026 19:19:59 +0400 Subject: [PATCH] docs(events,events-amqp): recovery backoff tuning, publisher shutdown, 1.2.0 currency MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - events-amqp: precise recovery table semantics (symmetric jitter, cap overshoot, per-series maxRetries); new Tuning the reconnect backoff section with the verified full-jitter workaround + upstream links; 1.2.0 currency (failFastOnInitialSetupError row, onSetupFailed row + example, fail-fast notes in Connection Recovery) - events: Graceful Shutdown section — handler-drain semantics, await-before-stop publisher recipe, stopping-gate limitation (connectum#212), planned drainPublishTimeout (connectum#196) - en/api: TypeDoc regen (AmqpRecoveryOptions JSDoc + source-link shifts) Companion: Connectum-Framework/connectum#214 Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_015Sp8jYmRXyvUmsErDuNJvt --- .../events-amqp/classes/AmqpAdapterError.md | 4 +- .../classes/AmqpConnectionError.md | 4 +- .../classes/AmqpPublishNackError.md | 4 +- .../classes/AmqpPublishTimeoutError.md | 4 +- .../classes/AmqpSerializationError.md | 4 +- .../events-amqp/classes/AmqpTopologyError.md | 4 +- .../classes/AmqpUnroutableError.md | 6 +- .../events-amqp/functions/AmqpAdapter.md | 2 +- .../types/interfaces/AmqpAdapterOptions.md | 61 ++++++++++++++++--- .../interfaces/AmqpBindingDeclaration.md | 12 ++-- .../types/interfaces/AmqpConsumerOptions.md | 6 +- .../interfaces/AmqpExchangeDeclaration.md | 12 ++-- .../types/interfaces/AmqpExchangeOptions.md | 6 +- .../interfaces/AmqpLifecycleCallbacks.md | 54 ++++++++++++++-- .../types/interfaces/AmqpPublisherOptions.md | 10 +-- .../types/interfaces/AmqpQueueDeclaration.md | 12 ++-- .../types/interfaces/AmqpQueueOptions.md | 12 ++-- .../types/interfaces/AmqpQueueOverride.md | 8 +-- .../types/interfaces/AmqpRecoveryOptions.md | 39 +++++++++--- .../interfaces/AmqpSerializationOptions.md | 8 +-- .../types/interfaces/AmqpTopology.md | 8 +-- .../types/type-aliases/AmqpTopologyMode.md | 2 +- .../types/variables/AmqpTopologyMode.md | 2 +- en/packages/events-amqp.md | 38 ++++++++++-- en/packages/events.md | 24 ++++++++ 25 files changed, 256 insertions(+), 90 deletions(-) diff --git a/en/api/@connectum/events-amqp/classes/AmqpAdapterError.md b/en/api/@connectum/events-amqp/classes/AmqpAdapterError.md index c2838bd6..addf318f 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpAdapterError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpAdapterError.md @@ -2,7 +2,7 @@ # Class: AmqpAdapterError -Defined in: [packages/events-amqp/src/errors.ts:13](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L13) +Defined in: [packages/events-amqp/src/errors.ts:30](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L30) Base class for all AMQP adapter errors. @@ -25,7 +25,7 @@ Base class for all AMQP adapter errors. > **new AmqpAdapterError**(`message`, `options?`): `AmqpAdapterError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpConnectionError.md b/en/api/@connectum/events-amqp/classes/AmqpConnectionError.md index 5c1a9056..840542dd 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpConnectionError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpConnectionError.md @@ -2,7 +2,7 @@ # Class: AmqpConnectionError -Defined in: [packages/events-amqp/src/errors.ts:25](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L25) +Defined in: [packages/events-amqp/src/errors.ts:42](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L42) Connection is absent, lost, or recovery is in progress / exhausted. Publishes during a disconnected window fail fast with this error; @@ -18,7 +18,7 @@ in-flight confirms are rejected with it on connection loss. > **new AmqpConnectionError**(`message`, `options?`): `AmqpConnectionError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpPublishNackError.md b/en/api/@connectum/events-amqp/classes/AmqpPublishNackError.md index b0012121..7c383555 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpPublishNackError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpPublishNackError.md @@ -2,7 +2,7 @@ # Class: AmqpPublishNackError -Defined in: [packages/events-amqp/src/errors.ts:41](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L41) +Defined in: [packages/events-amqp/src/errors.ts:58](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L58) The broker negatively acknowledged (nacked) a published message. @@ -16,7 +16,7 @@ The broker negatively acknowledged (nacked) a published message. > **new AmqpPublishNackError**(`message`, `options?`): `AmqpPublishNackError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpPublishTimeoutError.md b/en/api/@connectum/events-amqp/classes/AmqpPublishTimeoutError.md index 450fd475..7c86a0c4 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpPublishTimeoutError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpPublishTimeoutError.md @@ -2,7 +2,7 @@ # Class: AmqpPublishTimeoutError -Defined in: [packages/events-amqp/src/errors.ts:48](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L48) +Defined in: [packages/events-amqp/src/errors.ts:65](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L65) No broker outcome (ack/nack/return/connection loss) arrived within `publishTimeoutMs`. The message state is UNKNOWN — it may or may not @@ -18,7 +18,7 @@ have been routed; an at-least-once producer should republish. > **new AmqpPublishTimeoutError**(`message`, `options?`): `AmqpPublishTimeoutError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpSerializationError.md b/en/api/@connectum/events-amqp/classes/AmqpSerializationError.md index afaa04fd..5543d64e 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpSerializationError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpSerializationError.md @@ -2,7 +2,7 @@ # Class: AmqpSerializationError -Defined in: [packages/events-amqp/src/errors.ts:58](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L58) +Defined in: [packages/events-amqp/src/errors.ts:75](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L75) Payload encoding/decoding failed in a custom serialization hook. @@ -16,7 +16,7 @@ Payload encoding/decoding failed in a custom serialization hook. > **new AmqpSerializationError**(`message`, `options?`): `AmqpSerializationError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpTopologyError.md b/en/api/@connectum/events-amqp/classes/AmqpTopologyError.md index 76be3959..e06f1031 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpTopologyError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpTopologyError.md @@ -2,7 +2,7 @@ # Class: AmqpTopologyError -Defined in: [packages/events-amqp/src/errors.ts:55](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L55) +Defined in: [packages/events-amqp/src/errors.ts:72](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L72) Topology declaration or verification failed: missing exchange/queue in `check`/`skip` mode, or a conflicting redeclare (PRECONDITION_FAILED) in @@ -18,7 +18,7 @@ Topology declaration or verification failed: missing exchange/queue in > **new AmqpTopologyError**(`message`, `options?`): `AmqpTopologyError` -Defined in: [packages/events-amqp/src/errors.ts:14](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L14) +Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) #### Parameters diff --git a/en/api/@connectum/events-amqp/classes/AmqpUnroutableError.md b/en/api/@connectum/events-amqp/classes/AmqpUnroutableError.md index b91b6149..3b261aef 100644 --- a/en/api/@connectum/events-amqp/classes/AmqpUnroutableError.md +++ b/en/api/@connectum/events-amqp/classes/AmqpUnroutableError.md @@ -2,7 +2,7 @@ # Class: AmqpUnroutableError -Defined in: [packages/events-amqp/src/errors.ts:31](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L31) +Defined in: [packages/events-amqp/src/errors.ts:48](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L48) The broker returned a `mandatory` message as unroutable (`basic.return`): no queue is bound for the routing key. @@ -17,7 +17,7 @@ The broker returned a `mandatory` message as unroutable > **new AmqpUnroutableError**(`message`, `routingKey`): `AmqpUnroutableError` -Defined in: [packages/events-amqp/src/errors.ts:34](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L34) +Defined in: [packages/events-amqp/src/errors.ts:51](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L51) #### Parameters @@ -79,7 +79,7 @@ Defined in: node\_modules/.pnpm/typescript@5.9.3/node\_modules/typescript/lib/li > `readonly` **routingKey**: `string` -Defined in: [packages/events-amqp/src/errors.ts:32](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L32) +Defined in: [packages/events-amqp/src/errors.ts:49](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/errors.ts#L49) *** diff --git a/en/api/@connectum/events-amqp/functions/AmqpAdapter.md b/en/api/@connectum/events-amqp/functions/AmqpAdapter.md index 500cc4e9..64b34759 100644 --- a/en/api/@connectum/events-amqp/functions/AmqpAdapter.md +++ b/en/api/@connectum/events-amqp/functions/AmqpAdapter.md @@ -4,7 +4,7 @@ > **AmqpAdapter**(`options`): `EventAdapter` -Defined in: [packages/events-amqp/src/AmqpAdapter.ts:153](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/AmqpAdapter.ts#L153) +Defined in: [packages/events-amqp/src/AmqpAdapter.ts:254](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/AmqpAdapter.ts#L254) Create an AMQP/RabbitMQ adapter for @connectum/events. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpAdapterOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpAdapterOptions.md index 56a516fe..569ff822 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpAdapterOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpAdapterOptions.md @@ -60,11 +60,47 @@ Exchange type. *** +### failFastOnInitialSetupError? + +> `readonly` `optional` **failFastOnInitialSetupError?**: `boolean` + +Defined in: [packages/events-amqp/src/types.ts:149](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L149) + +Fail fast on a DETERMINISTIC setup/topology error on the FIRST connect, +instead of entering amqplib's infinite recovery loop. + +amqplib's opt-in recovery resolves `connect()` only after its setup hook +succeeds, and rejects only once `maxRetries` is exhausted (default +`Infinity`). A permanent topology error on the first connect under the +default recovery therefore HANGS `connect()` forever, with no thrown error +and — because the lifecycle listeners attach only after that never-returning +await — no callback. When this flag is `true` (and recovery is enabled), the +adapter first validates topology against a throwaway non-recovering +connection; a topology error rejects `connect()` with the typed +`AmqpTopologyError` / `AmqpConnectionError`. + +Only deterministic setup/topology errors fail fast. A transient +broker-unreachable at startup is NOT a fail-fast condition — it falls +through to normal recovery (block-until-broker). SUBSEQUENT reconnects +always keep infinite-recovery behavior. + +No-op with `recovery: false` (that path already fails fast on setup). +Enabling this (or supplying [AmqpLifecycleCallbacks.onSetupFailed](AmqpLifecycleCallbacks.md#onsetupfailed)) +adds one extra short-lived connection at startup for the validation probe. + +#### Default + +```ts +false +``` + +*** + ### lifecycle? > `readonly` `optional` **lifecycle?**: [`AmqpLifecycleCallbacks`](AmqpLifecycleCallbacks.md) -Defined in: [packages/events-amqp/src/types.ts:117](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L117) +Defined in: [packages/events-amqp/src/types.ts:155](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L155) Connection lifecycle callbacks. Connection errors are surfaced here — not just logged. @@ -85,7 +121,7 @@ Publisher options. > `readonly` `optional` **publishTimeoutMs?**: `number` -Defined in: [packages/events-amqp/src/types.ts:127](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L127) +Defined in: [packages/events-amqp/src/types.ts:165](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L165) Per-publish broker-outcome deadline in milliseconds. A publish whose ack/nack/return/connection-loss outcome does not arrive in time @@ -114,7 +150,7 @@ Default queue assertion options. > `readonly` `optional` **queueOverrides?**: `Record`\<`string`, [`AmqpQueueOverride`](AmqpQueueOverride.md)\> -Defined in: [packages/events-amqp/src/types.ts:97](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L97) +Defined in: [packages/events-amqp/src/types.ts:101](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L101) Map a consumer group name to an externally-named queue. @@ -128,7 +164,7 @@ lets a subscription attach to a queue from an external contract > `readonly` `optional` **recovery?**: `boolean` \| [`AmqpRecoveryOptions`](AmqpRecoveryOptions.md) -Defined in: [packages/events-amqp/src/types.ts:111](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L111) +Defined in: [packages/events-amqp/src/types.ts:122](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L122) Automatic connection recovery (delegated to amqplib's opt-in recovery). Enabled by default; pass `false` to restore @@ -139,6 +175,13 @@ topology (per `topologyMode`), and replays active subscriptions. In-flight publishes at the moment of a connection loss reject with `AmqpConnectionError`. +`maxRetries` governs BOTH the initial connect and steady-state recovery +(counter reset on success); under the default `Infinity`, `connect()` +blocks until the broker is reachable rather than failing fast (see +[AmqpAdapterOptions.failFastOnInitialSetupError](#failfastoninitialsetuperror) to fail fast on a +deterministic startup misconfiguration). See [AmqpRecoveryOptions](AmqpRecoveryOptions.md) +for the retry-budget scope and jitter/`maxDelay` overshoot. + #### Default ```ts @@ -190,12 +233,16 @@ exchange-to-exchange. > `readonly` `optional` **topologyMode?**: [`AmqpTopologyMode`](../type-aliases/AmqpTopologyMode.md) -Defined in: [packages/events-amqp/src/types.ts:88](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L88) +Defined in: [packages/events-amqp/src/types.ts:92](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L92) How topology is established: - `"assert"` (default) — declare idempotently (assertExchange/assertQueue/bind); -- `"check"` — existence-only verification (checkExchange/checkQueue), fail - fast with AmqpTopologyError on missing objects. AMQP offers no passive +- `"check"` — existence-only verification (checkExchange/checkQueue). A + missing object raises AmqpTopologyError, which fails `connect()` fast + ONLY with `recovery: false` or `failFastOnInitialSetupError: true`; under + the default recovery a first-connect check failure otherwise enters the + (infinite) recovery loop and is surfaced via `onSetupFailed` / + `onReconnecting` rather than rejecting `connect()`. AMQP offers no passive introspection: argument equivalence and binding presence are NOT verifiable in this mode (a conflicting redeclare elsewhere is PRECONDITION_FAILED 406); diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpBindingDeclaration.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpBindingDeclaration.md index 8b5b8317..421d3dd5 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpBindingDeclaration.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpBindingDeclaration.md @@ -2,7 +2,7 @@ # Interface: AmqpBindingDeclaration -Defined in: [packages/events-amqp/src/types.ts:187](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L187) +Defined in: [packages/events-amqp/src/types.ts:225](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L225) ## Properties @@ -10,7 +10,7 @@ Defined in: [packages/events-amqp/src/types.ts:187](https://github.com/Connectum > `readonly` `optional` **arguments?**: `Record`\<`string`, `unknown`\> -Defined in: [packages/events-amqp/src/types.ts:195](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L195) +Defined in: [packages/events-amqp/src/types.ts:233](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L233) *** @@ -18,7 +18,7 @@ Defined in: [packages/events-amqp/src/types.ts:195](https://github.com/Connectum > `readonly` `optional` **exchange?**: `string` -Defined in: [packages/events-amqp/src/types.ts:191](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L191) +Defined in: [packages/events-amqp/src/types.ts:229](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L229) Destination exchange name (exchange-to-exchange binding). @@ -28,7 +28,7 @@ Destination exchange name (exchange-to-exchange binding). > `readonly` `optional` **queue?**: `string` -Defined in: [packages/events-amqp/src/types.ts:189](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L189) +Defined in: [packages/events-amqp/src/types.ts:227](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L227) Destination queue name (queue binding) — mutually exclusive with `exchange`. @@ -38,7 +38,7 @@ Destination queue name (queue binding) — mutually exclusive with `exchange`. > `readonly` **routingKey**: `string` -Defined in: [packages/events-amqp/src/types.ts:194](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L194) +Defined in: [packages/events-amqp/src/types.ts:232](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L232) *** @@ -46,6 +46,6 @@ Defined in: [packages/events-amqp/src/types.ts:194](https://github.com/Connectum > `readonly` **source**: `string` -Defined in: [packages/events-amqp/src/types.ts:193](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L193) +Defined in: [packages/events-amqp/src/types.ts:231](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L231) Source exchange. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpConsumerOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpConsumerOptions.md index 2cc843ac..7ee136a4 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpConsumerOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpConsumerOptions.md @@ -2,7 +2,7 @@ # Interface: AmqpConsumerOptions -Defined in: [packages/events-amqp/src/types.ts:284](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L284) +Defined in: [packages/events-amqp/src/types.ts:364](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L364) Consumer options. @@ -12,7 +12,7 @@ Consumer options. > `readonly` `optional` **exclusive?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:298](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L298) +Defined in: [packages/events-amqp/src/types.ts:378](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L378) Whether the consumer is exclusive to this connection. @@ -28,7 +28,7 @@ false > `readonly` `optional` **prefetch?**: `number` -Defined in: [packages/events-amqp/src/types.ts:291](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L291) +Defined in: [packages/events-amqp/src/types.ts:371](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L371) Prefetch count (QoS) — how many unacknowledged messages a consumer can have at a time. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeDeclaration.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeDeclaration.md index f5ad4cf8..cc5f5df6 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeDeclaration.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeDeclaration.md @@ -2,7 +2,7 @@ # Interface: AmqpExchangeDeclaration -Defined in: [packages/events-amqp/src/types.ts:169](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L169) +Defined in: [packages/events-amqp/src/types.ts:207](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L207) ## Properties @@ -10,7 +10,7 @@ Defined in: [packages/events-amqp/src/types.ts:169](https://github.com/Connectum > `readonly` `optional` **arguments?**: `Record`\<`string`, `unknown`\> -Defined in: [packages/events-amqp/src/types.ts:175](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L175) +Defined in: [packages/events-amqp/src/types.ts:213](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L213) Raw AMQP arguments passthrough. @@ -20,7 +20,7 @@ Raw AMQP arguments passthrough. > `readonly` `optional` **autoDelete?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:173](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L173) +Defined in: [packages/events-amqp/src/types.ts:211](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L211) *** @@ -28,7 +28,7 @@ Defined in: [packages/events-amqp/src/types.ts:173](https://github.com/Connectum > `readonly` `optional` **durable?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:172](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L172) +Defined in: [packages/events-amqp/src/types.ts:210](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L210) *** @@ -36,7 +36,7 @@ Defined in: [packages/events-amqp/src/types.ts:172](https://github.com/Connectum > `readonly` **name**: `string` -Defined in: [packages/events-amqp/src/types.ts:170](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L170) +Defined in: [packages/events-amqp/src/types.ts:208](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L208) *** @@ -44,4 +44,4 @@ Defined in: [packages/events-amqp/src/types.ts:170](https://github.com/Connectum > `readonly` **type**: `"headers"` \| `"topic"` \| `"direct"` \| `"fanout"` -Defined in: [packages/events-amqp/src/types.ts:171](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L171) +Defined in: [packages/events-amqp/src/types.ts:209](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L209) diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeOptions.md index ec8043d7..81f5b3f5 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpExchangeOptions.md @@ -2,7 +2,7 @@ # Interface: AmqpExchangeOptions -Defined in: [packages/events-amqp/src/types.ts:233](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L233) +Defined in: [packages/events-amqp/src/types.ts:313](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L313) Exchange assertion options. @@ -12,7 +12,7 @@ Exchange assertion options. > `readonly` `optional` **autoDelete?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:246](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L246) +Defined in: [packages/events-amqp/src/types.ts:326](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L326) Whether the exchange is deleted when the last queue unbinds. @@ -28,7 +28,7 @@ false > `readonly` `optional` **durable?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:239](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L239) +Defined in: [packages/events-amqp/src/types.ts:319](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L319) Whether the exchange should survive broker restarts. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpLifecycleCallbacks.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpLifecycleCallbacks.md index 1aa82d93..875cf2f9 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpLifecycleCallbacks.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpLifecycleCallbacks.md @@ -2,7 +2,7 @@ # Interface: AmqpLifecycleCallbacks -Defined in: [packages/events-amqp/src/types.ts:223](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L223) +Defined in: [packages/events-amqp/src/types.ts:284](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L284) Connection lifecycle callbacks. @@ -12,7 +12,7 @@ Connection lifecycle callbacks. > `readonly` `optional` **onConnected?**: () => `void` -Defined in: [packages/events-amqp/src/types.ts:224](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L224) +Defined in: [packages/events-amqp/src/types.ts:285](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L285) #### Returns @@ -24,7 +24,7 @@ Defined in: [packages/events-amqp/src/types.ts:224](https://github.com/Connectum > `readonly` `optional` **onDisconnected?**: (`cause`) => `void` -Defined in: [packages/events-amqp/src/types.ts:225](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L225) +Defined in: [packages/events-amqp/src/types.ts:286](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L286) #### Parameters @@ -42,7 +42,7 @@ Defined in: [packages/events-amqp/src/types.ts:225](https://github.com/Connectum > `readonly` `optional` **onReconnectFailed?**: (`cause`) => `void` -Defined in: [packages/events-amqp/src/types.ts:227](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L227) +Defined in: [packages/events-amqp/src/types.ts:294](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L294) #### Parameters @@ -60,7 +60,12 @@ Defined in: [packages/events-amqp/src/types.ts:227](https://github.com/Connectum > `readonly` `optional` **onReconnecting?**: (`info`) => `void` -Defined in: [packages/events-amqp/src/types.ts:226](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L226) +Defined in: [packages/events-amqp/src/types.ts:293](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L293) + +A reconnect attempt has been scheduled. Fires exactly ONCE per scheduled +retry (amqplib's `reconnect-scheduled`). A failed attempt that also emits +`connect-failed` does NOT double-invoke this; the terminal, retries-exhausted +case is reported via [onReconnectFailed](#onreconnectfailed), not here. #### Parameters @@ -81,3 +86,42 @@ Defined in: [packages/events-amqp/src/types.ts:226](https://github.com/Connectum #### Returns `void` + +*** + +### onSetupFailed? + +> `readonly` `optional` **onSetupFailed?**: (`error`, `ctx`) => `void` + +Defined in: [packages/events-amqp/src/types.ts:307](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L307) + +A setup/topology failure occurred while (re)applying the declarative +topology — on the initial connect's validation probe (`ctx.initial: true`, +`ctx.attempt: 0`) and/or on a reconnect whose topology re-assert fails +(`ctx.initial: false`, `ctx.attempt` ≥ 1). + +This surfaces deterministic configuration drift (e.g. a missing queue in +`check` mode, or a `PRECONDITION_FAILED` redeclare) distinctly from a mere +broker outage, even when fail-fast is off. The initial-connect invocation +requires a startup validation probe, which runs when either this callback or +[AmqpAdapterOptions.failFastOnInitialSetupError](AmqpAdapterOptions.md#failfastoninitialsetuperror) is set. + +#### Parameters + +##### error + +`Error` + +##### ctx + +###### attempt + +`number` + +###### initial + +`boolean` + +#### Returns + +`void` diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpPublisherOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpPublisherOptions.md index 7d1ae9a4..1718634a 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpPublisherOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpPublisherOptions.md @@ -2,7 +2,7 @@ # Interface: AmqpPublisherOptions -Defined in: [packages/events-amqp/src/types.ts:304](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L304) +Defined in: [packages/events-amqp/src/types.ts:384](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L384) Publisher options. @@ -12,7 +12,7 @@ Publisher options. > `readonly` `optional` **correlationHeader?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:333](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L333) +Defined in: [packages/events-amqp/src/types.ts:413](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L413) How `basic.return` frames are correlated to publishes when `mandatory: true`. The return frame carries no deliveryTag, so: @@ -36,7 +36,7 @@ true > `readonly` `optional` **externalContract?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:362](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L362) +Defined in: [packages/events-amqp/src/types.ts:442](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L442) Publish against an EXTERNAL (non-EventBus) message contract: suppress the EventBus envelope so the wire frame carries ONLY contract-specified @@ -74,7 +74,7 @@ false > `readonly` `optional` **mandatory?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:318](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L318) +Defined in: [packages/events-amqp/src/types.ts:398](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L398) Whether the message should be returned if it cannot be routed. Unroutable messages reject the publish with `AmqpUnroutableError`. @@ -91,7 +91,7 @@ false > `readonly` `optional` **persistent?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:310](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L310) +Defined in: [packages/events-amqp/src/types.ts:390](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L390) Whether messages should be persisted to disk (deliveryMode=2). diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueDeclaration.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueDeclaration.md index 1384d9a6..776ebd7a 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueDeclaration.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueDeclaration.md @@ -2,7 +2,7 @@ # Interface: AmqpQueueDeclaration -Defined in: [packages/events-amqp/src/types.ts:178](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L178) +Defined in: [packages/events-amqp/src/types.ts:216](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L216) ## Properties @@ -10,7 +10,7 @@ Defined in: [packages/events-amqp/src/types.ts:178](https://github.com/Connectum > `readonly` `optional` **arguments?**: `Record`\<`string`, `unknown`\> -Defined in: [packages/events-amqp/src/types.ts:184](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L184) +Defined in: [packages/events-amqp/src/types.ts:222](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L222) Raw AMQP arguments passthrough (e.g. x-dead-letter-exchange). @@ -20,7 +20,7 @@ Raw AMQP arguments passthrough (e.g. x-dead-letter-exchange). > `readonly` `optional` **autoDelete?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:181](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L181) +Defined in: [packages/events-amqp/src/types.ts:219](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L219) *** @@ -28,7 +28,7 @@ Defined in: [packages/events-amqp/src/types.ts:181](https://github.com/Connectum > `readonly` `optional` **durable?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:180](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L180) +Defined in: [packages/events-amqp/src/types.ts:218](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L218) *** @@ -36,7 +36,7 @@ Defined in: [packages/events-amqp/src/types.ts:180](https://github.com/Connectum > `readonly` `optional` **exclusive?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:182](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L182) +Defined in: [packages/events-amqp/src/types.ts:220](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L220) *** @@ -44,4 +44,4 @@ Defined in: [packages/events-amqp/src/types.ts:182](https://github.com/Connectum > `readonly` **name**: `string` -Defined in: [packages/events-amqp/src/types.ts:179](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L179) +Defined in: [packages/events-amqp/src/types.ts:217](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L217) diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOptions.md index ee3e175e..264e9ec0 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOptions.md @@ -2,7 +2,7 @@ # Interface: AmqpQueueOptions -Defined in: [packages/events-amqp/src/types.ts:252](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L252) +Defined in: [packages/events-amqp/src/types.ts:332](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L332) Queue assertion options. @@ -12,7 +12,7 @@ Queue assertion options. > `readonly` `optional` **deadLetterExchange?**: `string` -Defined in: [packages/events-amqp/src/types.ts:273](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L273) +Defined in: [packages/events-amqp/src/types.ts:353](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L353) Dead letter exchange name for rejected messages. @@ -22,7 +22,7 @@ Dead letter exchange name for rejected messages. > `readonly` `optional` **deadLetterRoutingKey?**: `string` -Defined in: [packages/events-amqp/src/types.ts:278](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L278) +Defined in: [packages/events-amqp/src/types.ts:358](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L358) Dead letter routing key for rejected messages. @@ -32,7 +32,7 @@ Dead letter routing key for rejected messages. > `readonly` `optional` **durable?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:258](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L258) +Defined in: [packages/events-amqp/src/types.ts:338](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L338) Whether the queue should survive broker restarts. @@ -48,7 +48,7 @@ true > `readonly` `optional` **maxLength?**: `number` -Defined in: [packages/events-amqp/src/types.ts:268](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L268) +Defined in: [packages/events-amqp/src/types.ts:348](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L348) Maximum number of messages in the queue. @@ -58,6 +58,6 @@ Maximum number of messages in the queue. > `readonly` `optional` **messageTtl?**: `number` -Defined in: [packages/events-amqp/src/types.ts:263](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L263) +Defined in: [packages/events-amqp/src/types.ts:343](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L343) Per-message TTL in milliseconds. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOverride.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOverride.md index 0a283dff..7a13ac66 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOverride.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpQueueOverride.md @@ -2,7 +2,7 @@ # Interface: AmqpQueueOverride -Defined in: [packages/events-amqp/src/types.ts:199](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L199) +Defined in: [packages/events-amqp/src/types.ts:237](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L237) External queue override for a consumer group. @@ -12,7 +12,7 @@ External queue override for a consumer group. > `readonly` `optional` **arguments?**: `Record`\<`string`, `unknown`\> -Defined in: [packages/events-amqp/src/types.ts:203](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L203) +Defined in: [packages/events-amqp/src/types.ts:241](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L241) Raw AMQP arguments used when asserting the queue (assert mode only). @@ -22,7 +22,7 @@ Raw AMQP arguments used when asserting the queue (assert mode only). > `readonly` `optional` **durable?**: `boolean` -Defined in: [packages/events-amqp/src/types.ts:205](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L205) +Defined in: [packages/events-amqp/src/types.ts:243](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L243) #### Default @@ -36,6 +36,6 @@ true > `readonly` **queue**: `string` -Defined in: [packages/events-amqp/src/types.ts:201](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L201) +Defined in: [packages/events-amqp/src/types.ts:239](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L239) Externally-defined queue name to consume from. diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpRecoveryOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpRecoveryOptions.md index 110653a5..1e9adf2b 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpRecoveryOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpRecoveryOptions.md @@ -2,17 +2,38 @@ # Interface: AmqpRecoveryOptions -Defined in: [packages/events-amqp/src/types.ts:209](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L209) +Defined in: [packages/events-amqp/src/types.ts:270](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L270) Recovery knobs (passed through to amqplib's opt-in recovery). +`maxRetries` governs BOTH the initial connect and every subsequent recovery +series, with the counter reset on each success — so a finite value chosen only +to bound startup also caps steady-state recovery and makes the adapter brittle +(N consecutive transient failures in any single series stop it permanently). +The effective reconnect delay is symmetric jitter around the exponential +base — uniform in `[base × (1 − jitter), base × (1 + jitter)]` with +`base = min(maxDelay, initialDelay × factor^(attempt − 1))`. The cap applies +BEFORE jitter, so the wait can overshoot `maxDelay` (~20% at the default +jitter, up to ~2x at `jitter: 1`). + +Full jitter with a hard cap is expressible today: set `jitter: 1` and halve +`initialDelay`/`maxDelay` — the delay becomes uniform in `[0, intended cap]` +(verified against amqplib 2.0.1's internal formula; re-verify on upgrades). + +Bounding the initial connect independently from steady-state recovery, and a +pluggable backoff hook, are tracked as future options — see +[https://github.com/Connectum-Framework/connectum/issues/198](https://github.com/Connectum-Framework/connectum/issues/198) and +[https://github.com/Connectum-Framework/connectum/issues/199](https://github.com/Connectum-Framework/connectum/issues/199) +(upstream: [https://github.com/amqp-node/amqplib/issues/856](https://github.com/amqp-node/amqplib/issues/856) and +[https://github.com/amqp-node/amqplib/issues/855](https://github.com/amqp-node/amqplib/issues/855)). + ## Properties ### factor? > `readonly` `optional` **factor?**: `number` -Defined in: [packages/events-amqp/src/types.ts:215](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L215) +Defined in: [packages/events-amqp/src/types.ts:276](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L276) #### Default @@ -26,7 +47,7 @@ Defined in: [packages/events-amqp/src/types.ts:215](https://github.com/Connectum > `readonly` `optional` **initialDelay?**: `number` -Defined in: [packages/events-amqp/src/types.ts:211](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L211) +Defined in: [packages/events-amqp/src/types.ts:272](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L272) #### Default @@ -40,9 +61,9 @@ Defined in: [packages/events-amqp/src/types.ts:211](https://github.com/Connectum > `readonly` `optional` **jitter?**: `number` -Defined in: [packages/events-amqp/src/types.ts:217](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L217) +Defined in: [packages/events-amqp/src/types.ts:278](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L278) -0..1 +Symmetric jitter factor (0..1): the delay is uniform in `[base × (1 − jitter), base × (1 + jitter)]`. #### Default @@ -56,7 +77,9 @@ Defined in: [packages/events-amqp/src/types.ts:217](https://github.com/Connectum > `readonly` `optional` **maxDelay?**: `number` -Defined in: [packages/events-amqp/src/types.ts:213](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L213) +Defined in: [packages/events-amqp/src/types.ts:274](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L274) + +Base delay cap in ms; jitter is applied on top of the capped base, so the effective wait can exceed it. #### Default @@ -70,7 +93,9 @@ Defined in: [packages/events-amqp/src/types.ts:213](https://github.com/Connectum > `readonly` `optional` **maxRetries?**: `number` -Defined in: [packages/events-amqp/src/types.ts:219](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L219) +Defined in: [packages/events-amqp/src/types.ts:280](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L280) + +Attempts per series (initial connect and each recovery series); resets on success. #### Default diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpSerializationOptions.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpSerializationOptions.md index ec04a3b0..ea1ec817 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpSerializationOptions.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpSerializationOptions.md @@ -2,7 +2,7 @@ # Interface: AmqpSerializationOptions -Defined in: [packages/events-amqp/src/types.ts:140](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L140) +Defined in: [packages/events-amqp/src/types.ts:178](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L178) Serialization metadata and optional wire transcoding. @@ -12,7 +12,7 @@ Serialization metadata and optional wire transcoding. > `readonly` `optional` **contentType?**: `string` -Defined in: [packages/events-amqp/src/types.ts:146](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L146) +Defined in: [packages/events-amqp/src/types.ts:184](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L184) AMQP `contentType` message property. @@ -28,7 +28,7 @@ AMQP `contentType` message property. > `readonly` `optional` **decode?**: (`content`) => `Uint8Array` -Defined in: [packages/events-amqp/src/types.ts:159](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L159) +Defined in: [packages/events-amqp/src/types.ts:197](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L197) Transform the incoming wire body before it reaches the event handler. Failures nack the message (requeue per consumer policy). @@ -49,7 +49,7 @@ Failures nack the message (requeue per consumer policy). > `readonly` `optional` **encode?**: (`payload`) => `Uint8Array` -Defined in: [packages/events-amqp/src/types.ts:153](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L153) +Defined in: [packages/events-amqp/src/types.ts:191](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L191) Transform the outgoing wire body. Receives the payload bytes the EventBus (or the application) produced. Failures reject the publish diff --git a/en/api/@connectum/events-amqp/types/interfaces/AmqpTopology.md b/en/api/@connectum/events-amqp/types/interfaces/AmqpTopology.md index 03952af0..ecef8687 100644 --- a/en/api/@connectum/events-amqp/types/interfaces/AmqpTopology.md +++ b/en/api/@connectum/events-amqp/types/interfaces/AmqpTopology.md @@ -2,7 +2,7 @@ # Interface: AmqpTopology -Defined in: [packages/events-amqp/src/types.ts:163](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L163) +Defined in: [packages/events-amqp/src/types.ts:201](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L201) Declarative topology. @@ -12,7 +12,7 @@ Declarative topology. > `readonly` `optional` **bindings?**: readonly [`AmqpBindingDeclaration`](AmqpBindingDeclaration.md)[] -Defined in: [packages/events-amqp/src/types.ts:166](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L166) +Defined in: [packages/events-amqp/src/types.ts:204](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L204) *** @@ -20,7 +20,7 @@ Defined in: [packages/events-amqp/src/types.ts:166](https://github.com/Connectum > `readonly` `optional` **exchanges?**: readonly [`AmqpExchangeDeclaration`](AmqpExchangeDeclaration.md)[] -Defined in: [packages/events-amqp/src/types.ts:164](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L164) +Defined in: [packages/events-amqp/src/types.ts:202](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L202) *** @@ -28,4 +28,4 @@ Defined in: [packages/events-amqp/src/types.ts:164](https://github.com/Connectum > `readonly` `optional` **queues?**: readonly [`AmqpQueueDeclaration`](AmqpQueueDeclaration.md)[] -Defined in: [packages/events-amqp/src/types.ts:165](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L165) +Defined in: [packages/events-amqp/src/types.ts:203](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L203) diff --git a/en/api/@connectum/events-amqp/types/type-aliases/AmqpTopologyMode.md b/en/api/@connectum/events-amqp/types/type-aliases/AmqpTopologyMode.md index a514153e..7a453bd6 100644 --- a/en/api/@connectum/events-amqp/types/type-aliases/AmqpTopologyMode.md +++ b/en/api/@connectum/events-amqp/types/type-aliases/AmqpTopologyMode.md @@ -4,6 +4,6 @@ > **AmqpTopologyMode** = *typeof* [`AmqpTopologyMode`](../variables/AmqpTopologyMode.md)\[keyof *typeof* [`AmqpTopologyMode`](../variables/AmqpTopologyMode.md)\] -Defined in: [packages/events-amqp/src/types.ts:131](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L131) +Defined in: [packages/events-amqp/src/types.ts:169](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L169) Topology establishment mode. diff --git a/en/api/@connectum/events-amqp/types/variables/AmqpTopologyMode.md b/en/api/@connectum/events-amqp/types/variables/AmqpTopologyMode.md index 15744862..cc64e7f5 100644 --- a/en/api/@connectum/events-amqp/types/variables/AmqpTopologyMode.md +++ b/en/api/@connectum/events-amqp/types/variables/AmqpTopologyMode.md @@ -4,7 +4,7 @@ > `const` **AmqpTopologyMode**: `object` -Defined in: [packages/events-amqp/src/types.ts:131](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L131) +Defined in: [packages/events-amqp/src/types.ts:169](https://github.com/Connectum-Framework/connectum/blob/main/packages/events-amqp/src/types.ts#L169) Topology establishment mode. diff --git a/en/packages/events-amqp.md b/en/packages/events-amqp.md index 8aceed52..75ada38c 100644 --- a/en/packages/events-amqp.md +++ b/en/packages/events-amqp.md @@ -83,6 +83,7 @@ Pass the result to `createEventBus({ adapter })`. | `topologyMode` | `"assert" \| "check" \| "skip"` | `"assert"` | How topology is established | | `queueOverrides` | `Record` | `undefined` | Map a consumer group to an externally named queue | | `recovery` | `boolean \| AmqpRecoveryOptions` | `true` | Automatic connection recovery (amqplib native); `false` disables | +| `failFastOnInitialSetupError` | `boolean` | `false` | Reject `connect()` with the typed `AmqpTopologyError` on a deterministic setup/topology error at the **first** connect, instead of hanging in infinite recovery. Transient broker-unreachable still blocks-and-retries. Available since 1.2.0 | | `lifecycle` | `AmqpLifecycleCallbacks` | `undefined` | Connection lifecycle callbacks | | `publishTimeoutMs` | `number` | `30000` | Per-publish broker-outcome deadline in milliseconds | @@ -182,10 +183,10 @@ const adapter = AmqpAdapter({ | Option | Type | Default | Description | |--------|------|---------|-------------| | `initialDelay` | `number` | `100` | First reconnect delay in ms | -| `maxDelay` | `number` | `30000` | Delay cap in ms | +| `maxDelay` | `number` | `30000` | Base delay cap in ms. Jitter is applied on top of the capped base, so the effective wait can exceed this (~20% at the default jitter, up to ~2x at `jitter: 1`) | | `factor` | `number` | `2` | Exponential backoff factor | -| `jitter` | `number` | `0.2` | Randomization factor (0..1) | -| `maxRetries` | `number` | `Infinity` | Give up after this many attempts | +| `jitter` | `number` | `0.2` | Symmetric jitter factor (0..1): the delay is drawn uniformly from `[base × (1 − jitter), base × (1 + jitter)]` | +| `maxRetries` | `number` | `Infinity` | Attempts per series before giving up. Governs **both** the initial connect and each later recovery series; the counter resets on every success | ### `AmqpLifecycleCallbacks` @@ -193,8 +194,9 @@ const adapter = AmqpAdapter({ |----------|-----------|------------| | `onConnected` | `() => void` | Connection established (initial connect and after each recovery) | | `onDisconnected` | `(cause: Error) => void` | Connection lost | -| `onReconnecting` | `(info: { attempt: number; delay: number; error: Error }) => void` | A reconnect attempt is scheduled | +| `onReconnecting` | `(info: { attempt: number; delay: number; error: Error }) => void` | A reconnect attempt is scheduled (fires exactly once per scheduled retry) | | `onReconnectFailed` | `(cause: Error) => void` | Recovery exhausted (`maxRetries` reached) | +| `onSetupFailed` | `(error: Error, ctx: { initial: boolean; attempt: number }) => void` | Topology/setup failed on the initial validation probe (`initial: true`) or a reconnect re-assert (`initial: false`). The initial-connect call requires the startup probe, which runs when this callback or `failFastOnInitialSetupError` is set. Available since 1.2.0 | Connection errors are surfaced through these callbacks -- never console-only. @@ -245,6 +247,7 @@ const adapter = AmqpAdapter({ onDisconnected: (cause) => console.error('AMQP disconnected', cause), onReconnecting: ({ attempt, delay }) => console.warn(`Reconnect #${attempt} in ${delay}ms`), onReconnectFailed: (cause) => console.error('AMQP recovery exhausted', cause), + onSetupFailed: (error, { initial }) => console.error('AMQP topology/setup failed', { initial }, error), }, publishTimeoutMs: 30_000, }); @@ -393,10 +396,33 @@ On every successful (re)connect the adapter: Connection behavior: -- **With recovery enabled**, `connect()` retries with backoff until the broker becomes reachable -- convenient for `docker-compose` startup ordering where the broker may not be up yet. +- **With recovery enabled**, `connect()` retries with backoff until the broker becomes reachable -- convenient for `docker-compose` startup ordering where the broker may not be up yet. Under the default `maxRetries: Infinity`, `connect()` blocks rather than failing fast; set `failFastOnInitialSetupError: true` to reject `connect()` with a typed `AmqpTopologyError` on a **permanent** setup/topology error at startup while still recovering from transient broker outages. +- **`maxRetries` scope.** The retry budget governs **both** the initial connect and every later recovery series, with the counter reset on each success. A finite value chosen only to bound startup therefore also caps steady-state recovery: a transient blip of that many consecutive failures in any single series permanently stops recovery. +- **Reconnect delay.** The effective delay is symmetric jitter around the exponential base -- uniform in `[base × (1 − jitter), base × (1 + jitter)]` with `base = min(maxDelay, initialDelay × factor^(attempt − 1))`. The cap applies to the base **before** jitter, so the wait can overshoot `maxDelay` (~20% at the default `jitter: 0.2`, up to ~2x at `jitter: 1`). - **With `recovery: false`**, `connect()` rejects immediately if the broker is unreachable, and a lost connection is not restored. -Observe connection state through the `lifecycle` callbacks (`onConnected`, `onDisconnected`, `onReconnecting`, `onReconnectFailed`). +Observe connection state through the `lifecycle` callbacks (`onConnected`, `onDisconnected`, `onReconnecting`, `onReconnectFailed`, `onSetupFailed`). + +### Tuning the reconnect backoff + +The delay strategy is fixed inside amqplib -- a pluggable backoff hook and an independent initial-connect budget are proposed upstream ([amqp-node/amqplib#855](https://github.com/amqp-node/amqplib/issues/855), [amqp-node/amqplib#856](https://github.com/amqp-node/amqplib/issues/856); tracked in [connectum#199](https://github.com/Connectum-Framework/connectum/issues/199) and [connectum#198](https://github.com/Connectum-Framework/connectum/issues/198)). + +One shape that **is** expressible exactly with the current knobs is AWS-style **full jitter** with a hard cap -- the delay drawn uniformly from `[0, min(cap, schedule step)]`, never above the cap. Set `jitter: 1` and halve both `initialDelay` and `maxDelay`: + +```typescript +// Full jitter over an intended 500ms → 30s exponential schedule, hard-capped at 30s: +recovery: { + jitter: 1, // delay becomes uniform in [0, 2 × base] + initialDelay: 250, // half of the intended 500ms first step + maxDelay: 15_000, // half of the intended 30s cap +} +``` + +With `jitter: 1` the delay is uniform in `[0, 2 × base]`; halving the knobs makes `2 × base` trace the intended schedule, so the effective delay never exceeds the intended cap. + +::: warning Caveat +This leans on the exact internal delay formula of amqplib v2 (verified against 2.0.1: `base = min(maxDelay, initialDelay × factor^(attempt − 1))`, then a uniform offset of `± base × jitter`). It is precise today but is not a documented amqplib contract -- re-verify after amqplib upgrades. +::: ## Adapter Lifecycle diff --git a/en/packages/events.md b/en/packages/events.md index 56398a93..b49e3a71 100644 --- a/en/packages/events.md +++ b/en/packages/events.md @@ -315,6 +315,30 @@ const eventBus = createEventBus({ }); ``` +### Graceful Shutdown + +`stop()` closes subscriptions, waits for in-flight **consumer handlers** up to `drainTimeout` (default 30s), force-aborts the rest via `AbortSignal`, then disconnects the adapter. Set `drainTimeout: 0` for immediate abort. + +In-flight `publish()` promises are **not** tracked by the bus, and `drainTimeout` does not cover them. An at-least-once producer must settle its publishes before stopping: + +```typescript +// Track publishes you must not lose: +const pending = new Set>(); + +const p = bus.publish(OrderCreatedSchema, order); +pending.add(p); +p.catch(() => {}).finally(() => pending.delete(p)); + +// On shutdown — settle them BEFORE stop(): +await Promise.allSettled([...pending]); +await bus.stop(); +``` + +Two related boundaries: + +- **Publishing from a draining handler is rejected** — once `stop()` begins, `publish()` throws, including from handlers that are still draining. Relay topologies (consume → transform → publish) lose the in-flight tail at shutdown; the design discussion is tracked in [connectum#212](https://github.com/Connectum-Framework/connectum/issues/212). +- **An opt-in symmetric publish drain** (`drainPublishTimeout`) is planned — tracked in [connectum#196](https://github.com/Connectum-Framework/connectum/issues/196). + ## Exports Summary | Export | Description |