Skip to content

@amqp-contract/core


@amqp-contract/core

Classes

AmqpClient

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

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
  • Channel creation with JSON serialization enabled by default

All operations return AsyncResult<T, TechnicalError> for consistent error handling.

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', 'routingKey', { data: 'value' });

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

Constructors

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

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

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 channel with JSON serialization enabled
Parameters
ParameterTypeDescription
contractContractDefinitionThe contract definition specifying the AMQP topology
optionsAmqpClientOptionsClient configuration options
Returns

AmqpClient

Methods

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

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

Acknowledge a message.

Parameters
ParameterTypeDefault valueDescription
msgConsumeMessageundefinedThe message to acknowledge
allUpTobooleanfalseIf true, acknowledge all messages up to and including this one
Returns

void

addSetup()
ts
addSetup(setup): void;

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

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, TechnicalError>;

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

Cancel a consumer by its consumer tag.

Parameters
ParameterType
consumerTagstring
Returns

AsyncResult<void, TechnicalError>

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

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

Close the channel and release the connection reference.

This will:

  • Close the channel wrapper
  • Decrease the reference count on the shared connection
  • Close the connection if this was the last client using it

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

Returns

AsyncResult<void, TechnicalError>

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

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

Start consuming messages from a queue.

If options.prefetch is set, a per-consumer prefetch count is applied via channel.prefetch(count, false) registered as a setup function on the channel wrapper before the underlying consume call. Registering it via addSetup ensures the prefetch is reapplied automatically on channel reconnect; using global=false scopes it to subsequent consumers on the channel (RabbitMQ semantics — opposite of intuition: false is per- consumer, true is channel-wide).

prefetch is stripped from the options handed to channelWrapper.consume because it is not a valid amqplib Options.Consume field — leaving it in would just travel as a no-op key-value pair on the consume frame.

Parameters
ParameterType
queuestring
callbackConsumeCallback
options?ConsumerOptions
Returns

AsyncResult<string, TechnicalError>

AsyncResult resolving to the consumer tag.

getConnection()
ts
getConnection(): IAmqpConnectionManager;

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

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, 
   allUpTo?, 
   requeue?): void;

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

Negative acknowledge a message.

Parameters
ParameterTypeDefault valueDescription
msgConsumeMessageundefinedThe message to nack
allUpTobooleanfalseIf true, nack all messages up to and including this one
requeuebooleantrueIf true, requeue the message(s)
Returns

void

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

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

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(
   exchange, 
   routingKey, 
   content, 
   options?): AsyncResult<boolean, TechnicalError>;

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

Publish a message to an exchange.

Parameters
ParameterType
exchangestring
routingKeystring
contentunknown
options?Publish
Returns

AsyncResult<boolean, TechnicalError>

AsyncResult resolving to true if the message was sent, false if the channel buffer is full.

sendToQueue()
ts
sendToQueue(
   queue, 
   content, 
   options?): AsyncResult<boolean, TechnicalError>;

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

Publish a message directly to a queue.

Parameters
ParameterType
queuestring
contentunknown
options?Publish
Returns

AsyncResult<boolean, TechnicalError>

AsyncResult resolving to true if the message was sent, false if the channel buffer is full.

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

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

Wait for the channel to be connected and ready.

If connectTimeoutMs was provided in the constructor options, the returned AsyncResult resolves to Err(TechnicalError) 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 error path to release the connection — waitForConnect does not do this automatically. The typed factories handle this cleanup for you.

Returns

AsyncResult<void, TechnicalError>


MessageValidationError

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

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:45

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@4.1.0/node_modules/unthrown/dist/index.d.mts:1456
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:43
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:42
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:92

A typed, contract-declared RPC error — the business-failure channel of an RPC, as opposed to the transport failures modeled by TechnicalError.

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 matchTags; 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:102

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@4.1.0/node_modules/unthrown/dist/index.d.mts:1456
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:99
datareadonlyTDataTaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).data-packages/core/src/errors.ts:100
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:17

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

This includes AMQP connection failures, channel issues, validation failures, and other runtime errors. This error is shared across core, worker, and client packages.

Built on unthrown's TaggedError, so it carries a _tag of "@amqp-contract/TechnicalError" for exhaustive dispatch via matchTags. The tag is namespaced to avoid colliding with other libraries' tags in a shared matchTags; the human-facing Error.name is kept bare ("TechnicalError"). Remains a real Error (and a modeled error — it lives in the E channel of a Result, never the Defect channel).

Extends

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

Constructors

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

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

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@4.1.0/node_modules/unthrown/dist/index.d.mts:1456
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:81

Options for creating an AMQP client.

Properties

PropertyTypeDescriptionDefined in
channelOptions?Partial<CreateChannelOpts>Optional channel configuration options.packages/core/src/amqp-client.ts:84
connectionOptions?AmqpConnectionManagerOptionsOptional connection configuration (heartbeat, reconnect settings, etc.).packages/core/src/amqp-client.ts:83
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:85
urlsConnectionUrl[]AMQP broker URL(s). Multiple URLs provide failover support.packages/core/src/amqp-client.ts:82

ConsumeCallback

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

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

Callback type for consuming messages.

Parameters

ParameterType
msgConsumeMessage | null

Returns

void | Promise<void>


ConsumerOptions

ts
type ConsumerOptions = Options.Consume & object;

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

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

prefetch is intercepted by AmqpClient.consume: it is stripped from the options handed to the underlying channelWrapper.consume(...) call (since amqplib's Options.Consume does not include it) and applied via channel.prefetch(count, false) registered through addSetup before the consume so the value is in effect when the consumer starts and is reapplied automatically on channel reconnect.

Type Declaration

NameTypeDescriptionDefined in
prefetch?numberPer-consumer prefetch count. Applied before channel.consume(...).packages/core/src/amqp-client.ts:118

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

PublishOptions

ts
type PublishOptions = Options.Publish;

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

Publish options for AmqpClient.publish / AmqpClient.sendToQueue.

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.)


TelemetryProvider

ts
type TelemetryProvider = object;

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

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:71
getConsumeLatencyHistogram() => Histogram | undefinedGet a histogram for consume/process latency. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:83
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:90
getPublishCounter() => Counter | undefinedGet a counter for messages published. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:65
getPublishLatencyHistogram() => Histogram | undefinedGet a histogram for publish latency. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:77
getTracer() => Tracer | undefinedGet a tracer instance for creating spans. Returns undefined if OpenTelemetry is not available.packages/core/src/telemetry.ts:59

Variables

DEFAULT_CONNECT_TIMEOUT_MS

ts
const DEFAULT_CONNECT_TIMEOUT_MS: 30000 = 30_000;

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

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 — Infinity and other non-finite values are also coerced to "no timeout" because Node's setTimeout clamps large delays to ~24.8 days and silently fires near-immediately on Infinity.


defaultTelemetryProvider

ts
const defaultTelemetryProvider: TelemetryProvider;

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

Default telemetry provider that uses OpenTelemetry API if available.


MessagingSemanticConventions

ts
const MessagingSemanticConventions: object;

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

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:37
AMQP_PUBLISHER_NAME"amqp.publisher.name""amqp.publisher.name"packages/core/src/telemetry.ts:36
ERROR_TYPE"error.type""error.type"packages/core/src/telemetry.ts:40
MESSAGING_DESTINATION"messaging.destination.name""messaging.destination.name"packages/core/src/telemetry.ts:29
MESSAGING_DESTINATION_KIND"messaging.destination.kind""messaging.destination.kind"packages/core/src/telemetry.ts:30
MESSAGING_DESTINATION_KIND_EXCHANGE"exchange""exchange"packages/core/src/telemetry.ts:44
MESSAGING_DESTINATION_KIND_QUEUE"queue""queue"packages/core/src/telemetry.ts:45
MESSAGING_OPERATION"messaging.operation""messaging.operation"packages/core/src/telemetry.ts:31
MESSAGING_OPERATION_PROCESS"process""process"packages/core/src/telemetry.ts:47
MESSAGING_OPERATION_PUBLISH"publish""publish"packages/core/src/telemetry.ts:46
MESSAGING_RABBITMQ_MESSAGE_DELIVERY_TAG"messaging.rabbitmq.message.delivery_tag""messaging.rabbitmq.message.delivery_tag"packages/core/src/telemetry.ts:35
MESSAGING_RABBITMQ_ROUTING_KEY"messaging.rabbitmq.destination.routing_key""messaging.rabbitmq.destination.routing_key"packages/core/src/telemetry.ts:34
MESSAGING_SYSTEM"messaging.system""messaging.system"packages/core/src/telemetry.ts:28
MESSAGING_SYSTEM_RABBITMQ"rabbitmq""rabbitmq"packages/core/src/telemetry.ts:43

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:74

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

_getConnectionCountForTesting()

ts
function _getConnectionCountForTesting(): number;

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

Returns

number

Deprecated

Renamed to _internal_getConnectionCount per the org _internal_ convention.


_resetConnectionsForTesting()

ts
function _resetConnectionsForTesting(): Promise<void>;

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

Returns

Promise<void>

Deprecated

Renamed to _internal_resetConnections per the org _internal_ convention.


_resetTelemetryCacheForTesting()

ts
function _resetTelemetryCacheForTesting(): void;

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

Returns

void

Deprecated

Renamed to _internal_resetTelemetryCache per the org _internal_ convention.


endSpanError()

ts
function endSpanError(span, error): void;

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

End a span with error status.

Parameters

ParameterType
spanSpan | undefined
errorError

Returns

void


endSpanSuccess()

ts
function endSpanSuccess(span): void;

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

End a span with success status.

Parameters

ParameterType
spanSpan | undefined

Returns

void


isRpcError()

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

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

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 matchTags on the _tag.

Parameters

ParameterType
errorunknown

Returns

error is RpcError<string, unknown>


recordConsumeMetric()

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

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

Record a consume metric.

Parameters

ParameterType
providerTelemetryProvider
queueNamestring
consumerNamestring
successboolean
durationMsnumber

Returns

void


recordLateRpcReply()

ts
function recordLateRpcReply(provider, reason): void;

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

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:341

Record a publish metric.

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:143

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<E>(buffer, errorFn): Result<unknown, E>;

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

Parse a Buffer as JSON, mapping any JSON.parse exception to the caller-supplied error type.

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).

Type Parameters

Type ParameterDescription
EThe error type produced by errorFn.

Parameters

ParameterTypeDescription
bufferBufferThe raw message body to parse.
errorFn(raw) => ECallback invoked with the underlying JSON.parse error.

Returns

Result<unknown, E>

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

Example

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

setupAmqpTopology()

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

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

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)
  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:277

Create a span for a consume/process operation. Returns undefined if OpenTelemetry is not available.

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:242

Create a span for a publish operation. Returns undefined if OpenTelemetry is not available.

Parameters

ParameterType
providerTelemetryProvider
exchangeNamestring
routingKeystring | undefined
attributes?Attributes

Returns

Span | undefined

Released under the MIT License.