Skip to content

@amqp-contract/core


@amqp-contract/core

Classes

AmqpClient

Defined in: packages/core/src/amqp-client.ts:242

AMQP client that manages connections and channels with automatic topology setup.

This class handles:

  • Connection management with automatic reconnection via amqp-connection-manager
  • Connection pooling and sharing across instances with the same URLs
  • Automatic AMQP topology setup (exchanges, queues, bindings) from contract
  • Content encoding: non-Buffer payloads are JSON-encoded at publish time, Buffers go on the wire byte-for-byte

All operations return AsyncResult<T, never>: infrastructure failures are unexpected, so they surface through the Defect channel (with a TechnicalError as the defect's cause for logging), never as a modeled Err.

Example

typescript
const client = new AmqpClient(contract, {
  urls: ['amqp://localhost'],
  connectionOptions: { heartbeatIntervalInSeconds: 30 }
});

// Wait for connection (AsyncResult is thenable)
await client.waitForConnect();

// Publish a message
const result = await client.publish(
  { exchange: 'exchange', routingKey: 'routingKey' },
  { data: 'value' },
);

// Close when done
await client.close().get();

Constructors

Constructor
ts
new AmqpClient(contract, options): AmqpClient;

Defined in: packages/core/src/amqp-client.ts:275

Create a new AMQP client instance.

The client will automatically:

  • Get or create a shared connection using the singleton pattern
  • Set up AMQP topology (exchanges, queues, bindings) from the contract
  • Create a confirm channel that encodes content at publish time (JSON for plain values, byte-for-byte for Buffers)
Parameters
ParameterTypeDescription
contractContractDefinitionThe contract definition specifying the AMQP topology
optionsAmqpClientOptionsClient configuration options
Returns

AmqpClient

Accessors

currentChannelEpoch
Get Signature
ts
get currentChannelEpoch(): number;

Defined in: packages/core/src/amqp-client.ts:568

The current channel epoch — bumped on every channel 'connect' (initial connect included). Consumers stamp deliveries with this value and pass it back via ack / nack so a settle can never target a tag from a previous channel incarnation.

Returns

number

Methods

ack()
ts
ack(msg, options?): void;

Defined in: packages/core/src/amqp-client.ts:607

Acknowledge a message.

Parameters
ParameterTypeDescription
msgConsumeMessageThe message to acknowledge
options?{ allUpTo?: boolean; deliveryEpoch?: number; }Settle options: - allUpTo — if true, acknowledge all messages up to and including this one (defaults to false). - deliveryEpoch — pass the epoch captured when the message was delivered (currentChannelEpoch) to make the ack reconnect-safe: a stale epoch skips the ack (logged) instead of settling a foreign tag.
options.allUpTo?boolean-
options.deliveryEpoch?number-
Returns

void

addSetup()
ts
addSetup(setup): void;

Defined in: packages/core/src/amqp-client.ts:649

Add a setup function to be called when the channel is created or reconnected.

This is useful for setting up channel-level configuration like prefetch.

Parameters
ParameterTypeDescription
setup(channel) => void | Promise<void>The setup function to add
Returns

void

cancel()
ts
cancel(consumerTag): AsyncResult<void, never>;

Defined in: packages/core/src/amqp-client.ts:556

Cancel a consumer by its consumer tag.

Parameters
ParameterType
consumerTagstring
Returns

AsyncResult<void, never>

close()
ts
close(): AsyncResult<void, never>;

Defined in: packages/core/src/amqp-client.ts:682

Close the channel and release the connection lease.

This will:

  • Close the channel wrapper
  • Release this client's lease on the shared connection
  • Close the connection if this was the last client using it

Idempotent: a second close() returns the same in-flight (or settled) result instead of double-releasing the shared connection.

Both steps run regardless of each other's outcome; if both fail, the errors are wrapped in an AggregateError.

Returns

AsyncResult<void, never>

consume()
ts
consume(
   queue, 
   callback, 
   options?): AsyncResult<string, never>;

Defined in: packages/core/src/amqp-client.ts:518

Start consuming messages from a queue.

options.prefetch maps to amqp-connection-manager's native per-consumer prefetch: applied via basic.qos(count, global=false) immediately before this consumer's basic.consume, and re-applied the same way when the consumer is re-established after a reconnect — so the value never bleeds onto other consumers sharing the channel.

Parameters
ParameterType
queuestring
callbackConsumeCallback
options?AmqpConsumeOptions
Returns

AsyncResult<string, never>

AsyncResult resolving to the consumer tag.

getConnection()
ts
getConnection(): IAmqpConnectionManager;

Defined in: packages/core/src/amqp-client.ts:371

Get the underlying connection manager

This method exposes the AmqpConnectionManager instance that this client uses. The connection is automatically shared across all AmqpClient instances that use the same URLs and connection options.

Returns

IAmqpConnectionManager

The AmqpConnectionManager instance used by this client

nack()
ts
nack(msg, options?): void;

Defined in: packages/core/src/amqp-client.ts:628

Negative acknowledge a message.

Parameters
ParameterTypeDescription
msgConsumeMessageThe message to nack
options?{ allUpTo?: boolean; deliveryEpoch?: number; requeue?: boolean; }Settle options: - allUpTo — if true, nack all messages up to and including this one (defaults to false). - requeue — if true, requeue the message(s) (defaults to true). - deliveryEpoch — pass the epoch captured at delivery time to make the nack reconnect-safe (see ack).
options.allUpTo?boolean-
options.deliveryEpoch?number-
options.requeue?boolean-
Returns

void

on()
ts
on(event, listener): void;

Defined in: packages/core/src/amqp-client.ts:664

Register an event listener on the channel wrapper.

Available events:

  • 'connect': Emitted when the channel is (re)connected
  • 'close': Emitted when the channel is closed
  • 'error': Emitted when an error occurs
Parameters
ParameterTypeDescription
eventstringThe event name
listener(...args) => voidThe event listener
Returns

void

publish()
ts
publish(
   target, 
   content, 
   options?): AsyncResult<void, never>;

Defined in: packages/core/src/amqp-client.ts:454

Publish a message to an exchange.

Non-Buffer content is JSON-encoded; Buffers are published byte-for-byte.

A full channel write buffer (the wrapper's boolean false confirmation) surfaces as a Defect with a TechnicalError cause — like every other publish-side infrastructure failure. Callers never see the boolean.

Parameters
ParameterTypeDescription
target{ exchange: string; routingKey: string; }The exchange and routing key to publish to
target.exchangestring-
target.routingKeystring-
content?unknownThe message payload
options?PublishAMQP publish options
Returns

AsyncResult<void, never>

sendToQueue()
ts
sendToQueue(
   queue, 
   content, 
   options?): AsyncResult<void, never>;

Defined in: packages/core/src/amqp-client.ts:490

Publish a message directly to a queue.

Non-Buffer content is JSON-encoded; Buffers are published byte-for-byte.

A full channel write buffer surfaces as a Defect with a TechnicalError cause — see publish.

Parameters
ParameterType
queuestring
contentunknown
options?Publish
Returns

AsyncResult<void, never>

waitForConnect()
ts
waitForConnect(): AsyncResult<void, never>;

Defined in: packages/core/src/amqp-client.ts:390

Wait for the channel to be connected and ready.

If connectTimeoutMs was provided in the constructor options, the returned AsyncResult resolves to a Defect (a TechnicalError cause) once the timeout elapses. Without a timeout, this waits forever — amqp-connection-manager retries connections indefinitely and never errors on its own.

NOTE: When using AmqpClient directly (not via TypedAmqpClient / TypedAmqpWorker), the constructor has already incremented the pooled connection's reference count. Callers must invoke close() on the failure path to release the connection — waitForConnect does not do this automatically. The typed factories handle this cleanup for you.

Returns

AsyncResult<void, never>


MessageValidationError

Defined in: packages/core/src/errors.ts:46

Error thrown when message validation fails (payload or headers).

Used by both the client (publish-time payload validation) and the worker (consume-time payload and headers validation). Carries a _tag of "@amqp-contract/MessageValidationError" (namespaced to avoid collisions); the Error.name is kept bare ("MessageValidationError").

Param

source

The name of the publisher or consumer that triggered the validation

Param

issues

The validation issues from the Standard Schema validation

Extends

  • TaggedErrorInstance<"@amqp-contract/MessageValidationError", { issues: unknown; source: string; }>

Constructors

Constructor
ts
new MessageValidationError(source, issues): MessageValidationError;

Defined in: packages/core/src/errors.ts:52

Parameters
ParameterType
sourcestring
issuesunknown
Returns

MessageValidationError

Overrides
ts
TaggedError("@amqp-contract/MessageValidationError", {
  name: "MessageValidationError",
})<{
  source: string;
  issues: unknown;
}>.constructor

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/MessageValidationError"TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", })._tagnode_modules/.pnpm/unthrown@5.1.0/node_modules/unthrown/dist/index.d.mts:1941
cause?publicunknownTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).causenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24
issuesreadonlyunknownTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).issuespackages/core/src/errors.ts:50
messagepublicstringTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
sourcereadonlystringTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).sourcepackages/core/src/errors.ts:49
stack?publicstringTaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

RpcError

Defined in: packages/core/src/errors.ts:116

A typed, contract-declared RPC error — the business-failure channel of an RPC, as opposed to transport failures (which surface as a Defect with a TechnicalError cause).

Declared per-RPC via defineRpc(queue, { request, response, errors }), where each error code maps to a message definition validating the error's data payload. A worker handler surfaces one by returning Err(rpcError(code, data)); the worker validates data against the declared schema, publishes an error reply, and acks the request (business errors are not retried). The caller's client.call(...) resolves to Err(RpcError<code, data>) with data re-validated on arrival.

Carries a _tag of "@amqp-contract/RpcError" for exhaustive dispatch via the error matcher (result.match({ ok, defect, errCases: (matcher) => matcher.with(P.tag("@amqp-contract/RpcError"), …) })); the Error.name is kept bare ("RpcError"). Discriminate between codes on the code property.

Extends

  • TaggedErrorInstance<"@amqp-contract/RpcError", { code: string; data: unknown; }>

Type Parameters

Type ParameterDefault type
TCode extends stringstring
TDataunknown

Constructors

Constructor
ts
new RpcError<TCode, TData>(
   code, 
   data, 
   message?): RpcError<TCode, TData>;

Defined in: packages/core/src/errors.ts:126

Parameters
ParameterType
codeTCode
dataTData
message?string
Returns

RpcError<TCode, TData>

Overrides
ts
TaggedError(
  "@amqp-contract/RpcError",
  { name: "RpcError" },
)<{
  code: string;
  data: unknown;
}>.constructor

Properties

PropertyModifierTypeOverridesInherited fromDefined in
_tagreadonly"@amqp-contract/RpcError"-TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, )._tagnode_modules/.pnpm/unthrown@5.1.0/node_modules/unthrown/dist/index.d.mts:1941
cause?publicunknown-TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).causenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24
codereadonlyTCodeTaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).code-packages/core/src/errors.ts:123
datareadonlyTDataTaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).data-packages/core/src/errors.ts:124
messagepublicstring-TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstring-TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstring-TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

TechnicalError

Defined in: packages/core/src/errors.ts:24

Error for technical/runtime failures that cannot be prevented by TypeScript.

This includes AMQP connection failures, channel issues, compression/parse faults, and other unexpected runtime errors. Shared across core, worker, and client packages.

These failures are unexpected, so @amqp-contract surfaces them through unthrown's defect channel, not the modeled E channel: a TechnicalError instance is carried as the cause of a Defect (so its message/cause survive for logging), and is handled in the defect arm of result.match({ ok, errCases, defect }) — or via recoverDefect / tapDefect — never matched in errCases. It is deliberately absent from every operation's E (only anticipated domain failures live there).

Built on unthrown's TaggedError, so it carries a _tag of "@amqp-contract/TechnicalError" (namespaced to avoid colliding with other libraries' tags); the human-facing Error.name is kept bare ("TechnicalError"). Remains a real Error.

Extends

  • TaggedErrorInstance<"@amqp-contract/TechnicalError", { cause?: unknown; }>

Constructors

Constructor
ts
new TechnicalError(message, cause?): TechnicalError;

Defined in: packages/core/src/errors.ts:29

Parameters
ParameterType
messagestring
cause?unknown
Returns

TechnicalError

Overrides
ts
TaggedError("@amqp-contract/TechnicalError", {
  name: "TechnicalError",
})<{
  cause?: unknown;
}>.constructor

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/TechnicalError"TaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", })._tagnode_modules/.pnpm/unthrown@5.1.0/node_modules/unthrown/dist/index.d.mts:1941
cause?publicunknownMessageValidationError.causenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24
messagepublicstringTaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringTaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstringTaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

Type Aliases

AmqpClientOptions

ts
type AmqpClientOptions = object;

Defined in: packages/core/src/amqp-client.ts:133

Options for creating an AMQP client.

Properties

PropertyTypeDescriptionDefined in
channelOptions?Partial<CreateChannelOpts>Optional channel configuration options.packages/core/src/amqp-client.ts:136
connectionOptions?AmqpConnectionManagerOptionsOptional connection configuration (heartbeat, reconnect settings, etc.).packages/core/src/amqp-client.ts:135
connectTimeoutMs?number | nullMaximum time in ms to wait for the channel to become ready in waitForConnect. Defaults to DEFAULT_CONNECT_TIMEOUT_MS. Pass null to disable the timeout entirely (amqp-connection-manager will retry indefinitely).packages/core/src/amqp-client.ts:137
logger?LoggerOptional logger. Channel-level 'error' events (topology setup failures on connect/reconnect, publish-worker faults) are routed here — they are recoverable-by-reconnect conditions, never thrown.packages/core/src/amqp-client.ts:157
publishTimeoutMs?number | nullMaximum time in ms a publish may sit buffered waiting for the broker before its promise settles with a failure. Defaults to DEFAULT_PUBLISH_TIMEOUT_MS. Pass null to disable, restoring unbounded buffering. See the field's own doc comment for the precedence against channelOptions.publishTimeout.packages/core/src/amqp-client.ts:156
urlsConnectionUrl[]AMQP broker URL(s). Multiple URLs provide failover support.packages/core/src/amqp-client.ts:134

AmqpConsumeOptions

ts
type AmqpConsumeOptions = Omit<Options.Consume, "prefetch"> & object;

Defined in: packages/core/src/amqp-client.ts:193

Consume options that extend amqplib's Options.Consume with an optional per-consumer prefetch count.

Named AmqpConsumeOptions (not ConsumerOptions) so it never collides with the user-facing ConsumerOptions of @amqp-contract/worker.

prefetch maps to amqp-connection-manager's native per-consumer prefetch: it is applied via basic.qos(count, global=false) immediately before this consumer's basic.consume — and re-applied the same way when the consumer is re-established after a reconnect — so the value never bleeds onto other consumers sharing the channel.

Type Declaration

NameTypeDescriptionDefined in
prefetch?number | "unbounded"Per-consumer prefetch count, applied before channel.consume(...). Defaults to DEFAULT_PREFETCH. Pass "unbounded" to opt out and let the broker push the entire ready backlog — AMQP's original default, and a memory hazard on any queue that can build a backlog. "unbounded" rather than 0 because AMQP's 0 means unlimited, which reads at a call site as its opposite.packages/core/src/amqp-client.ts:204

AmqpPublishOptions

ts
type AmqpPublishOptions = Options.Publish;

Defined in: packages/core/src/amqp-client.ts:178

Publish options for AmqpClient.publish / AmqpClient.sendToQueue.

Named AmqpPublishOptions (not PublishOptions) so it never collides with the user-facing PublishOptions of @amqp-contract/client.

Currently a re-export of amqplib's Options.Publish. A previous version of this type also exposed a timeout field, but that field never had a meaningful AMQP-level effect in this codebase and has been removed to avoid suggesting behaviour we do not provide. (amqp-connection-manager's own publishTimeout channel option is unrelated and is configured at channel creation, not per-publish.)


ConnectionLease

ts
type ConnectionLease = object;

Defined in: packages/core/src/connection-manager.ts:10

A held reference to a pooled connection.

Properties

PropertyTypeDescriptionDefined in
connectionAmqpConnectionManagerThe shared connection this lease holds a reference to.packages/core/src/connection-manager.ts:12
release() => Promise<void>Release this lease. Idempotent — a double release() (e.g. a client whose close() runs on both a finally and an error path) is a no-op, so it can never underflow the pool's reference count and close a connection out from under another live client.packages/core/src/connection-manager.ts:19

ConsumeCallback

ts
type ConsumeCallback = (msg) => void | Promise<void>;

Defined in: packages/core/src/amqp-client.ts:163

Callback type for consuming messages.

Parameters

ParameterType
msgConsumeMessage | null

Returns

void | Promise<void>


Logger

ts
type Logger = object;

Defined in: packages/core/src/logger.ts:30

Logger interface for amqp-contract packages.

Provides a simple logging abstraction that can be implemented by users to integrate with their preferred logging framework.

Example

typescript
// Simple console logger implementation
const logger: Logger = {
  debug: (message, context) => console.debug(message, context),
  info: (message, context) => console.info(message, context),
  warn: (message, context) => console.warn(message, context),
  error: (message, context) => console.error(message, context),
};

Methods

debug()
ts
debug(message, context?): void;

Defined in: packages/core/src/logger.ts:36

Log debug level messages

Parameters
ParameterTypeDescription
messagestringThe log message
context?LoggerContextOptional context to include with the log
Returns

void

error()
ts
error(message, context?): void;

Defined in: packages/core/src/logger.ts:57

Log error level messages

Parameters
ParameterTypeDescription
messagestringThe log message
context?LoggerContextOptional context to include with the log
Returns

void

info()
ts
info(message, context?): void;

Defined in: packages/core/src/logger.ts:43

Log info level messages

Parameters
ParameterTypeDescription
messagestringThe log message
context?LoggerContextOptional context to include with the log
Returns

void

warn()
ts
warn(message, context?): void;

Defined in: packages/core/src/logger.ts:50

Log warning level messages

Parameters
ParameterTypeDescription
messagestringThe log message
context?LoggerContextOptional context to include with the log
Returns

void


LoggerContext

ts
type LoggerContext = Record<string, unknown> & object;

Defined in: packages/core/src/logger.ts:9

Context object for logger methods.

This type includes reserved keys that provide consistent naming for common logging context properties.

Type Declaration

NameTypeDefined in
error?unknownpackages/core/src/logger.ts:10

TelemetryProvider

ts
type TelemetryProvider = object;

Defined in: packages/core/src/telemetry.ts:55

Telemetry provider for AMQP operations. Uses lazy loading to gracefully handle cases where OpenTelemetry is not installed.

Properties

PropertyTypeDescriptionDefined in
getConsumeCounter() => Counter | undefinedGet a counter for messages consumed. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:72
getConsumeLatencyHistogram() => Histogram | undefinedGet a histogram for consume/process latency. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:84
getLateRpcReplyCounter() => Counter | undefinedGet a counter for RPC replies that arrive after the caller has gone away (timeout, cancellation, or unknown correlationId). Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:91
getPublishCounter() => Counter | undefinedGet a counter for messages published. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:66
getPublishLatencyHistogram() => Histogram | undefinedGet a histogram for publish latency. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:78
getTracer() => Tracer | undefinedGet a tracer instance for creating spans. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:60

Variables

DEFAULT_CONNECT_TIMEOUT_MS

ts
const DEFAULT_CONNECT_TIMEOUT_MS: 30000 = 30_000;

Defined in: packages/core/src/amqp-client.ts:70

Default time waitForConnect will wait for the broker before erroring out. Defaulting to a finite value (rather than waiting forever) means a fail-fast developer experience: a misconfigured URL, a down broker, or wrong credentials surface as an err within 30 seconds. Pass null explicitly to disable the timeout.


DEFAULT_PREFETCH

ts
const DEFAULT_PREFETCH: 10 = 10;

Defined in: packages/core/src/amqp-client.ts:79

Default per-consumer prefetch.

Bounds in-flight messages per consumer, which bounds both memory and the redelivery burst when a worker crashes. Throughput-bound consumers raise it explicitly; "unbounded" restores AMQP's unlimited behavior.


DEFAULT_PUBLISH_TIMEOUT_MS

ts
const DEFAULT_PUBLISH_TIMEOUT_MS: 30000 = 30_000;

Defined in: packages/core/src/amqp-client.ts:92

Default publishTimeout for the channel, in milliseconds.

Without a bound, publishes issued while the broker is unreachable buffer indefinitely and their promises never settle — a caller awaiting one waits forever. 30s is long enough that a brief reconnect does not fail healthy publishes, short enough that a real outage surfaces as an error.

Pass publishTimeoutMs: null to disable, matching the connectTimeoutMs convention.


defaultTelemetryProvider

ts
const defaultTelemetryProvider: TelemetryProvider;

Defined in: packages/core/src/telemetry.ts:230

Default telemetry provider that uses OpenTelemetry API if available.


MessagingSemanticConventions

ts
const MessagingSemanticConventions: object;

Defined in: packages/core/src/telemetry.ts:27

Semantic conventions for AMQP messaging following OpenTelemetry standards.

Type Declaration

NameTypeDefault valueDefined in
AMQP_CONSUMER_NAME"amqp.consumer.name""amqp.consumer.name"packages/core/src/telemetry.ts:38
AMQP_PUBLISHER_NAME"amqp.publisher.name""amqp.publisher.name"packages/core/src/telemetry.ts:37
ERROR_TYPE"error.type""error.type"packages/core/src/telemetry.ts:41
MESSAGING_DESTINATION"messaging.destination.name""messaging.destination.name"packages/core/src/telemetry.ts:30
MESSAGING_DESTINATION_KIND"messaging.destination.kind""messaging.destination.kind"packages/core/src/telemetry.ts:31
MESSAGING_DESTINATION_KIND_EXCHANGE"exchange""exchange"packages/core/src/telemetry.ts:45
MESSAGING_DESTINATION_KIND_QUEUE"queue""queue"packages/core/src/telemetry.ts:46
MESSAGING_OPERATION"messaging.operation""messaging.operation"packages/core/src/telemetry.ts:32
MESSAGING_OPERATION_PROCESS"process""process"packages/core/src/telemetry.ts:48
MESSAGING_OPERATION_PUBLISH"publish""publish"packages/core/src/telemetry.ts:47
MESSAGING_RABBITMQ_MESSAGE_DELIVERY_TAG"messaging.rabbitmq.message.delivery_tag""messaging.rabbitmq.message.delivery_tag"packages/core/src/telemetry.ts:36
MESSAGING_RABBITMQ_ROUTING_KEY"messaging.rabbitmq.destination.routing_key""messaging.rabbitmq.destination.routing_key"packages/core/src/telemetry.ts:35
MESSAGING_SYSTEM"messaging.system""messaging.system"packages/core/src/telemetry.ts:29
MESSAGING_SYSTEM_RABBITMQ"rabbitmq""rabbitmq"packages/core/src/telemetry.ts:44

See

https://opentelemetry.io/docs/specs/semconv/messaging/messaging-spans/


RPC_ERROR_CODE_HEADER

ts
const RPC_ERROR_CODE_HEADER: "x-amqp-contract-error-code" = "x-amqp-contract-error-code";

Defined in: packages/core/src/errors.ts:96

AMQP message header carrying the error code of a typed RPC error reply.

A reply message with this header is an error reply: its body is { message, data } where data conforms to the error's declared schema in the RPC's errors map. A reply without it is a regular success reply whose body is the response payload — so success replies are wire-compatible with contracts that declare no errors.

Functions

endSpanError()

ts
function endSpanError(span, error): void;

Defined in: packages/core/src/telemetry.ts:366

End a span with error status. Never throws.

Parameters

ParameterType
spanSpan | undefined
errorError

Returns

void


endSpanSuccess()

ts
function endSpanSuccess(span): void;

Defined in: packages/core/src/telemetry.ts:349

End a span with success status. Never throws.

Parameters

ParameterType
spanSpan | undefined

Returns

void


isMessageValidationError()

ts
function isMessageValidationError(error): error is MessageValidationError;

Defined in: packages/core/src/errors.ts:83

Type guard to check if an error is a MessageValidationError.

Parameters

ParameterType
errorunknown

Returns

error is MessageValidationError


isRpcError()

ts
function isRpcError(error): error is RpcError<string, unknown>;

Defined in: packages/core/src/errors.ts:139

Type guard to check if an error is an RpcError.

Narrowing to a specific code (and thus a typed data) is done on the code property after the guard, or via the error matcher on the _tag (matcher.with(P.tag("@amqp-contract/RpcError"), …)).

Parameters

ParameterType
errorunknown

Returns

error is RpcError<string, unknown>


isTechnicalError()

ts
function isTechnicalError(error): error is TechnicalError;

Defined in: packages/core/src/errors.ts:76

Type guard to check if an error is a TechnicalError — the cause carried by every infrastructure Defect this library produces.

Parameters

ParameterType
errorunknown

Returns

error is TechnicalError


recordConsumeMetric()

ts
function recordConsumeMetric(
   provider, 
   queueName, 
   consumerName, 
   success, 
   durationMs): void;

Defined in: packages/core/src/telemetry.ts:414

Record a consume metric. Never throws.

Parameters

ParameterType
providerTelemetryProvider
queueNamestring
consumerNamestring
successboolean
durationMsnumber

Returns

void


recordLateRpcReply()

ts
function recordLateRpcReply(provider, reason): void;

Defined in: packages/core/src/telemetry.ts:446

Record an RPC reply that arrived after the caller stopped waiting.

Parameters

ParameterTypeDescription
providerTelemetryProvider-
reason"unknown-correlation-id" | "missing-correlation-id"Why the reply was orphaned. "unknown-correlation-id" is the typical "caller already timed out" case; "missing-correlation-id" means the broker delivered a reply with no correlationId at all (a protocol violation by the responder).

Returns

void


recordPublishMetric()

ts
function recordPublishMetric(
   provider, 
   exchangeName, 
   routingKey, 
   success, 
   durationMs): void;

Defined in: packages/core/src/telemetry.ts:385

Record a publish metric. Never throws.

Parameters

ParameterType
providerTelemetryProvider
exchangeNamestring
routingKeystring | undefined
successboolean
durationMsnumber

Returns

void


rpcError()

ts
function rpcError<TCode, TData>(
   code, 
   data, 
   message?): RpcError<TCode, TData>;

Defined in: packages/core/src/errors.ts:168

Create an RpcError with less verbosity.

The code/data pair must match one of the entries declared in the RPC's errors map — the handler's return type enforces this at compile time, and the worker validates data against the declared schema at runtime before replying.

Type Parameters

Type Parameter
TCode extends string
TData

Parameters

ParameterTypeDescription
codeTCodeThe error code, as declared in the RPC's errors map
dataTDataThe error data, validated against the declared schema
message?stringOptional human-readable message (defaults to a generic one)

Returns

RpcError<TCode, TData>

Example

typescript
import { rpcError } from '@amqp-contract/worker';
import { ErrAsync } from 'unthrown';

const handler = ({ payload }) => {
  if (!orders.has(payload.orderId)) {
    return ErrAsync(rpcError('ORDER_NOT_FOUND', { orderId: payload.orderId }));
  }
  // ...
};

safeJsonParse()

ts
function safeJsonParse<R>(buffer, qualify): Result<unknown, Exclude<R, Defect>>;

Defined in: packages/core/src/parsing.ts:47

Parse a Buffer as JSON, triaging any JSON.parse exception through the caller-supplied qualify callback.

Use this in consume / reply paths where a parse failure is a typed value, not a thrown exception — the caller decides how to translate the raw error into a domain-level error (e.g. TechnicalError), or routes it to the defect channel via the injected defect helper (the full unthrown qualify signature, so no model-then-defect round-trip is ever needed).

Type Parameters

Type ParameterDescription
RWhat the qualify callback produces: a modeled error type, the defect marker, or a union of both. The defect marker is subtracted from the resulting error channel, exactly like fromThrowable.

Parameters

ParameterTypeDescription
bufferBufferThe raw message body to parse.
qualify(raw, defect) => R & NotThenable<R>Callback invoked with the underlying JSON.parse error and the injected defect helper; returns the modeled error or defect(cause).

Returns

Result<unknown, Exclude<R, Defect>>

A Result containing the parsed unknown value or the mapped error.

Examples

Modeled error

typescript
const parsed = safeJsonParse(
  msg.content,
  (error) => new TechnicalError("Failed to parse JSON", error),
);

Defect-channel routing

typescript
const parsed = safeJsonParse(
  msg.content,
  (error, defect) => defect(new TechnicalError("Failed to parse JSON", error)),
); // Result<unknown, never>

setupAmqpTopology()

ts
function setupAmqpTopology(channel, contract): Promise<void>;

Defined in: packages/core/src/setup.ts:30

Setup AMQP topology (exchanges, queues, and bindings) from a contract definition.

This function sets up the complete AMQP topology in the correct order:

  1. Assert all exchanges defined in the contract
  2. Validate dead letter exchanges are declared before referencing them
  3. Assert all queues with their configurations (including dead letter settings), plus the TTL-backoff wait queues derived from each queue's retry config (one per distinct backoff delay — see deriveTtlBackoffInfrastructure)
  4. Create all bindings (queue-to-exchange and exchange-to-exchange)

Parameters

ParameterTypeDescription
channelChannelThe AMQP channel to use for topology setup
contractContractDefinitionThe contract definition containing the topology specification

Returns

Promise<void>

Throws

If any exchanges, queues, or bindings fail to be created

Throws

If a queue references a dead letter exchange not declared in the contract

Example

typescript
const channel = await connection.createChannel();
await setupAmqpTopology(channel, contract);

startConsumeSpan()

ts
function startConsumeSpan(
   provider, 
   queueName, 
   consumerName, 
   attributes?): Span | undefined;

Defined in: packages/core/src/telemetry.ts:306

Create a span for a consume/process operation. Returns undefined if OpenTelemetry is not available. Never throws — a throwing provider is treated as "no telemetry".

Parameters

ParameterType
providerTelemetryProvider
queueNamestring
consumerNamestring
attributes?Attributes

Returns

Span | undefined


startPublishSpan()

ts
function startPublishSpan(
   provider, 
   exchangeName, 
   routingKey, 
   attributes?): Span | undefined;

Defined in: packages/core/src/telemetry.ts:259

Create a span for a publish operation. Returns undefined if OpenTelemetry is not available. Never throws — a throwing provider is treated as "no telemetry".

Parameters

ParameterType
providerTelemetryProvider
exchangeNamestring
routingKeystring | undefined
attributes?Attributes

Returns

Span | undefined


technicalDefect()

ts
function technicalDefect(error): Result<never, never>;

Defined in: packages/core/src/defect.ts:20

Mint a Defect-carrying Result from a TechnicalError, for the imperative sites (outside a combinator callback) that must surface an unexpected infrastructure failure through the defect channel. Uses the fromSafeThrowable boundary — the sanctioned way to route a throw to a Defect without a public defect constructor.

The shape is the sync Result<never, never>: it fits every channel (both T and E are never), so it can seed any pipeline. Call .toAsync() where an AsyncResult is expected.

Shared by @amqp-contract/client and @amqp-contract/worker (each used to hand-roll its own copy); exported so any layer can mint a defect the same way instead of re-deriving the boundary trick.

Parameters

ParameterType
errorTechnicalError

Returns

Result<never, never>

Released under the MIT License.