@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
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
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
| Parameter | Type | Description |
|---|---|---|
contract | ContractDefinition | The contract definition specifying the AMQP topology |
options | AmqpClientOptions | Client configuration options |
Returns
Methods
ack()
ack(msg, allUpTo?): void;Defined in: packages/core/src/amqp-client.ts:443
Acknowledge a message.
Parameters
| Parameter | Type | Default value | Description |
|---|---|---|---|
msg | ConsumeMessage | undefined | The message to acknowledge |
allUpTo | boolean | false | If true, acknowledge all messages up to and including this one |
Returns
void
addSetup()
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
| Parameter | Type | Description |
|---|---|---|
setup | (channel) => void | Promise<void> | The setup function to add |
Returns
void
cancel()
cancel(consumerTag): AsyncResult<void, TechnicalError>;Defined in: packages/core/src/amqp-client.ts:412
Cancel a consumer by its consumer tag.
Parameters
| Parameter | Type |
|---|---|
consumerTag | string |
Returns
AsyncResult<void, TechnicalError>
close()
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()
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
| Parameter | Type |
|---|---|
queue | string |
callback | ConsumeCallback |
options? | ConsumerOptions |
Returns
AsyncResult<string, TechnicalError>
AsyncResult resolving to the consumer tag.
getConnection()
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()
nack(
msg,
allUpTo?,
requeue?): void;Defined in: packages/core/src/amqp-client.ts:454
Negative acknowledge a message.
Parameters
| Parameter | Type | Default value | Description |
|---|---|---|---|
msg | ConsumeMessage | undefined | The message to nack |
allUpTo | boolean | false | If true, nack all messages up to and including this one |
requeue | boolean | true | If true, requeue the message(s) |
Returns
void
on()
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
| Parameter | Type | Description |
|---|---|---|
event | string | The event name |
listener | (...args) => void | The event listener |
Returns
void
publish()
publish(
exchange,
routingKey,
content,
options?): AsyncResult<boolean, TechnicalError>;Defined in: packages/core/src/amqp-client.ts:288
Publish a message to an exchange.
Parameters
| Parameter | Type |
|---|---|
exchange | string |
routingKey | string |
content | unknown |
options? | Publish |
Returns
AsyncResult<boolean, TechnicalError>
AsyncResult resolving to true if the message was sent, false if the channel buffer is full.
sendToQueue()
sendToQueue(
queue,
content,
options?): AsyncResult<boolean, TechnicalError>;Defined in: packages/core/src/amqp-client.ts:309
Publish a message directly to a queue.
Parameters
| Parameter | Type |
|---|---|
queue | string |
content | unknown |
options? | Publish |
Returns
AsyncResult<boolean, TechnicalError>
AsyncResult resolving to true if the message was sent, false if the channel buffer is full.
waitForConnect()
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
new MessageValidationError(source, issues): MessageValidationError;Defined in: packages/core/src/errors.ts:45
Parameters
| Parameter | Type |
|---|---|
source | string |
issues | unknown |
Returns
Overrides
TaggedError("@amqp-contract/MessageValidationError", {
name: "MessageValidationError",
})<{
source: string;
issues: unknown;
}>.constructorProperties
| Property | Modifier | Type | Inherited from | Defined in |
|---|---|---|---|---|
_tag | readonly | "@amqp-contract/MessageValidationError" | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", })._tag | node_modules/.pnpm/unthrown@4.1.0/node_modules/unthrown/dist/index.d.mts:1456 |
cause? | public | unknown | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).cause | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24 |
issues | readonly | unknown | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).issues | packages/core/src/errors.ts:43 |
message | public | string | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).message | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075 |
name | public | string | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).name | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074 |
source | readonly | string | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).source | packages/core/src/errors.ts:42 |
stack? | public | string | TaggedError("@amqp-contract/MessageValidationError", { name: "MessageValidationError", }).stack | node_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 Parameter | Default type |
|---|---|
TCode extends string | string |
TData | unknown |
Constructors
Constructor
new RpcError<TCode, TData>(
code,
data,
message?): RpcError<TCode, TData>;Defined in: packages/core/src/errors.ts:102
Parameters
| Parameter | Type |
|---|---|
code | TCode |
data | TData |
message? | string |
Returns
RpcError<TCode, TData>
Overrides
TaggedError(
"@amqp-contract/RpcError",
{ name: "RpcError" },
)<{
code: string;
data: unknown;
}>.constructorProperties
| Property | Modifier | Type | Overrides | Inherited from | Defined in |
|---|---|---|---|---|---|
_tag | readonly | "@amqp-contract/RpcError" | - | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, )._tag | node_modules/.pnpm/unthrown@4.1.0/node_modules/unthrown/dist/index.d.mts:1456 |
cause? | public | unknown | - | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).cause | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24 |
code | readonly | TCode | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).code | - | packages/core/src/errors.ts:99 |
data | readonly | TData | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).data | - | packages/core/src/errors.ts:100 |
message | public | string | - | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).message | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075 |
name | public | string | - | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).name | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074 |
stack? | public | string | - | TaggedError( "@amqp-contract/RpcError", { name: "RpcError" }, ).stack | node_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
new TechnicalError(message, cause?): TechnicalError;Defined in: packages/core/src/errors.ts:22
Parameters
| Parameter | Type |
|---|---|
message | string |
cause? | unknown |
Returns
Overrides
TaggedError("@amqp-contract/TechnicalError", {
name: "TechnicalError",
})<{
cause?: unknown;
}>.constructorProperties
| Property | Modifier | Type | Inherited from | Defined in |
|---|---|---|---|---|
_tag | readonly | "@amqp-contract/TechnicalError" | TaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", })._tag | node_modules/.pnpm/unthrown@4.1.0/node_modules/unthrown/dist/index.d.mts:1456 |
cause? | public | unknown | MessageValidationError.cause | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24 |
message | public | string | TaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).message | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075 |
name | public | string | TaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).name | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074 |
stack? | public | string | TaggedError("@amqp-contract/TechnicalError", { name: "TechnicalError", }).stack | node_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076 |
Type Aliases
AmqpClientOptions
type AmqpClientOptions = object;Defined in: packages/core/src/amqp-client.ts:81
Options for creating an AMQP client.
Properties
| Property | Type | Description | Defined in |
|---|---|---|---|
channelOptions? | Partial<CreateChannelOpts> | Optional channel configuration options. | packages/core/src/amqp-client.ts:84 |
connectionOptions? | AmqpConnectionManagerOptions | Optional connection configuration (heartbeat, reconnect settings, etc.). | packages/core/src/amqp-client.ts:83 |
connectTimeoutMs? | number | null | Maximum 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 |
urls | ConnectionUrl[] | AMQP broker URL(s). Multiple URLs provide failover support. | packages/core/src/amqp-client.ts:82 |
ConsumeCallback
type ConsumeCallback = (msg) => void | Promise<void>;Defined in: packages/core/src/amqp-client.ts:91
Callback type for consuming messages.
Parameters
| Parameter | Type |
|---|---|
msg | ConsumeMessage | null |
Returns
void | Promise<void>
ConsumerOptions
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
| Name | Type | Description | Defined in |
|---|---|---|---|
prefetch? | number | Per-consumer prefetch count. Applied before channel.consume(...). | packages/core/src/amqp-client.ts:118 |
Logger
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
// 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()
debug(message, context?): void;Defined in: packages/core/src/logger.ts:36
Log debug level messages
Parameters
| Parameter | Type | Description |
|---|---|---|
message | string | The log message |
context? | LoggerContext | Optional context to include with the log |
Returns
void
error()
error(message, context?): void;Defined in: packages/core/src/logger.ts:57
Log error level messages
Parameters
| Parameter | Type | Description |
|---|---|---|
message | string | The log message |
context? | LoggerContext | Optional context to include with the log |
Returns
void
info()
info(message, context?): void;Defined in: packages/core/src/logger.ts:43
Log info level messages
Parameters
| Parameter | Type | Description |
|---|---|---|
message | string | The log message |
context? | LoggerContext | Optional context to include with the log |
Returns
void
warn()
warn(message, context?): void;Defined in: packages/core/src/logger.ts:50
Log warning level messages
Parameters
| Parameter | Type | Description |
|---|---|---|
message | string | The log message |
context? | LoggerContext | Optional context to include with the log |
Returns
void
LoggerContext
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
| Name | Type | Defined in |
|---|---|---|
error? | unknown | packages/core/src/logger.ts:10 |
PublishOptions
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
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
| Property | Type | Description | Defined in |
|---|---|---|---|
getConsumeCounter | () => Counter | undefined | Get a counter for messages consumed. Returns undefined if OpenTelemetry is not available. | packages/core/src/telemetry.ts:71 |
getConsumeLatencyHistogram | () => Histogram | undefined | Get a histogram for consume/process latency. Returns undefined if OpenTelemetry is not available. | packages/core/src/telemetry.ts:83 |
getLateRpcReplyCounter | () => Counter | undefined | Get 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 | undefined | Get a counter for messages published. Returns undefined if OpenTelemetry is not available. | packages/core/src/telemetry.ts:65 |
getPublishLatencyHistogram | () => Histogram | undefined | Get a histogram for publish latency. Returns undefined if OpenTelemetry is not available. | packages/core/src/telemetry.ts:77 |
getTracer | () => Tracer | undefined | Get a tracer instance for creating spans. Returns undefined if OpenTelemetry is not available. | packages/core/src/telemetry.ts:59 |
Variables
DEFAULT_CONNECT_TIMEOUT_MS
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
const defaultTelemetryProvider: TelemetryProvider;Defined in: packages/core/src/telemetry.ts:229
Default telemetry provider that uses OpenTelemetry API if available.
MessagingSemanticConventions
const MessagingSemanticConventions: object;Defined in: packages/core/src/telemetry.ts:26
Semantic conventions for AMQP messaging following OpenTelemetry standards.
Type Declaration
| Name | Type | Default value | Defined 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
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()
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()
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()
function _resetTelemetryCacheForTesting(): void;Defined in: packages/core/src/telemetry.ts:429
Returns
void
Deprecated
Renamed to _internal_resetTelemetryCache per the org _internal_ convention.
endSpanError()
function endSpanError(span, error): void;Defined in: packages/core/src/telemetry.ts:324
End a span with error status.
Parameters
| Parameter | Type |
|---|---|
span | Span | undefined |
error | Error |
Returns
void
endSpanSuccess()
function endSpanSuccess(span): void;Defined in: packages/core/src/telemetry.ts:309
End a span with success status.
Parameters
| Parameter | Type |
|---|---|
span | Span | undefined |
Returns
void
isRpcError()
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
| Parameter | Type |
|---|---|
error | unknown |
Returns
error is RpcError<string, unknown>
recordConsumeMetric()
function recordConsumeMetric(
provider,
queueName,
consumerName,
success,
durationMs): void;Defined in: packages/core/src/telemetry.ts:368
Record a consume metric.
Parameters
| Parameter | Type |
|---|---|
provider | TelemetryProvider |
queueName | string |
consumerName | string |
success | boolean |
durationMs | number |
Returns
void
recordLateRpcReply()
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
| Parameter | Type | Description |
|---|---|---|
provider | TelemetryProvider | - |
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()
function recordPublishMetric(
provider,
exchangeName,
routingKey,
success,
durationMs): void;Defined in: packages/core/src/telemetry.ts:341
Record a publish metric.
Parameters
| Parameter | Type |
|---|---|
provider | TelemetryProvider |
exchangeName | string |
routingKey | string | undefined |
success | boolean |
durationMs | number |
Returns
void
rpcError()
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
| Parameter | Type | Description |
|---|---|---|
code | TCode | The error code, as declared in the RPC's errors map |
data | TData | The error data, validated against the declared schema |
message? | string | Optional human-readable message (defaults to a generic one) |
Returns
RpcError<TCode, TData>
Example
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()
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 Parameter | Description |
|---|---|
E | The error type produced by errorFn. |
Parameters
| Parameter | Type | Description |
|---|---|---|
buffer | Buffer | The raw message body to parse. |
errorFn | (raw) => E | Callback invoked with the underlying JSON.parse error. |
Returns
Result<unknown, E>
A Result containing the parsed unknown value or the mapped error.
Example
const parsed = safeJsonParse(
msg.content,
(error) => new TechnicalError("Failed to parse JSON", error),
);setupAmqpTopology()
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:
- Assert all exchanges defined in the contract
- Validate dead letter exchanges are declared before referencing them
- Assert all queues with their configurations (including dead letter settings)
- Create all bindings (queue-to-exchange and exchange-to-exchange)
Parameters
| Parameter | Type | Description |
|---|---|---|
channel | Channel | The AMQP channel to use for topology setup |
contract | ContractDefinition | The 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
const channel = await connection.createChannel();
await setupAmqpTopology(channel, contract);startConsumeSpan()
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
| Parameter | Type |
|---|---|
provider | TelemetryProvider |
queueName | string |
consumerName | string |
attributes? | Attributes |
Returns
Span | undefined
startPublishSpan()
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
| Parameter | Type |
|---|---|
provider | TelemetryProvider |
exchangeName | string |
routingKey | string | undefined |
attributes? | Attributes |
Returns
Span | undefined