Skip to content

@amqp-contract/worker


@amqp-contract/worker

Classes

MessageValidationError

Defined in: packages/core/dist/index.d.mts:454

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

  • MessageValidationError_base<{ issues: unknown; source: string; }>

Constructors

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

Defined in: packages/core/dist/index.d.mts:458

Parameters
ParameterType
sourcestring
issuesunknown
Returns

MessageValidationError

Overrides
ts
MessageValidationError_base<{
  source: string;
  issues: unknown;
}>.constructor

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/MessageValidationError"MessageValidationError_base._tagnode_modules/.pnpm/unthrown@5.1.0/node_modules/unthrown/dist/index.d.mts:1941
cause?publicunknownMessageValidationError_base.causenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24
issuesreadonlyunknownMessageValidationError_base.issuespackages/core/dist/index.d.mts:456
messagepublicstringMessageValidationError_base.messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringMessageValidationError_base.namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
sourcereadonlystringMessageValidationError_base.sourcepackages/core/dist/index.d.mts:455
stack?publicstringMessageValidationError_base.stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

NonRetryableError

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

Non-retryable errors - permanent failures that should not be retried Examples: invalid data, business rule violations, permanent external failures

Use this error type when retrying would not help - the message will be immediately sent to the dead letter queue (DLQ) if configured. Carries a namespaced _tag of "@amqp-contract/NonRetryableError"; the Error.name is kept bare ("NonRetryableError").

Extends

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

Constructors

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

Defined in: packages/worker/src/errors.ts:44

Parameters
ParameterType
messagestring
cause?unknown
Returns

NonRetryableError

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

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/NonRetryableError"TaggedError("@amqp-contract/NonRetryableError", { name: "NonRetryableError", })._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/NonRetryableError", { name: "NonRetryableError", }).messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringTaggedError("@amqp-contract/NonRetryableError", { name: "NonRetryableError", }).namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstringTaggedError("@amqp-contract/NonRetryableError", { name: "NonRetryableError", }).stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

RetryableError

Defined in: packages/worker/src/errors.ts:19

Retryable errors - transient failures that may succeed on retry Examples: network timeouts, rate limiting, temporary service unavailability

Use this error type when the operation might succeed if retried. The worker will apply exponential backoff and retry the message.

Built on unthrown's TaggedError, so it carries a namespaced _tag of "@amqp-contract/RetryableError" (to avoid colliding with other libraries' tags in a shared error matcher) for exhaustive dispatch via result.match({ ok, defect, errCases: (matcher) => matcher.with(P.tag("@amqp-contract/RetryableError"), …) }); the Error.name is kept bare ("RetryableError").

Extends

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

Constructors

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

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

Parameters
ParameterType
messagestring
cause?unknown
Returns

RetryableError

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

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/RetryableError"TaggedError("@amqp-contract/RetryableError", { name: "RetryableError", })._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/RetryableError", { name: "RetryableError", }).messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringTaggedError("@amqp-contract/RetryableError", { name: "RetryableError", }).namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstringTaggedError("@amqp-contract/RetryableError", { name: "RetryableError", }).stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

RpcError

Defined in: packages/core/dist/index.d.mts:498

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

  • RpcError_base<{ 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/dist/index.d.mts:504

Parameters
ParameterType
codeTCode
dataTData
message?string
Returns

RpcError<TCode, TData>

Overrides
ts
RpcError_base<{
  code: string;
  data: unknown;
}>.constructor

Properties

PropertyModifierTypeOverridesInherited fromDefined in
_tagreadonly"@amqp-contract/RpcError"-RpcError_base._tagnode_modules/.pnpm/unthrown@5.1.0/node_modules/unthrown/dist/index.d.mts:1941
cause?publicunknown-RpcError_base.causenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es2022.error.d.ts:24
codereadonlyTCodeRpcError_base.code-packages/core/dist/index.d.mts:502
datareadonlyTDataRpcError_base.data-packages/core/dist/index.d.mts:503
messagepublicstring-RpcError_base.messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstring-RpcError_base.namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstring-RpcError_base.stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

TechnicalError

Defined in: packages/core/dist/index.d.mts:437

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

  • TechnicalError_base<{ cause?: unknown; }>

Constructors

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

Defined in: packages/core/dist/index.d.mts:440

Parameters
ParameterType
messagestring
cause?unknown
Returns

TechnicalError

Overrides
ts
TechnicalError_base<{
  cause?: unknown;
}>.constructor

Properties

PropertyModifierTypeInherited fromDefined in
_tagreadonly"@amqp-contract/TechnicalError"TechnicalError_base._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
messagepublicstringTechnicalError_base.messagenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1075
namepublicstringTechnicalError_base.namenode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1074
stack?publicstringTechnicalError_base.stacknode_modules/.pnpm/typescript@6.0.3/node_modules/typescript/lib/lib.es5.d.ts:1076

TypedAmqpWorker

Defined in: packages/worker/src/worker.ts:331

Type-safe AMQP worker for consuming messages from RabbitMQ.

This class provides automatic message validation, connection management, and error handling for consuming messages based on a contract definition.

Example

typescript
import { TypedAmqpWorker } from '@amqp-contract/worker';
import { defineQueue, defineMessage, defineContract, defineConsumer } from '@amqp-contract/contract';
import { OkAsync } from 'unthrown';
import { z } from 'zod';

const orderQueue = defineQueue('order-processing');
const orderMessage = defineMessage(z.object({
  orderId: z.string(),
  amount: z.number()
}));

const contract = defineContract({
  consumers: {
    processOrder: defineConsumer(orderQueue, orderMessage)
  }
});

const worker = await TypedAmqpWorker.create({
  contract,
  handlers: {
    processOrder: ({ payload }) => {
      console.log('Processing order', payload.orderId);
      return OkAsync(undefined);
    },
  },
  urls: ['amqp://localhost'],
}).get();

// Close when done (drains in-flight handlers first)
await worker.close().get();

Type Parameters

Type ParameterDescription
TContract extends ContractDefinitionThe contract definition type

Methods

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

Defined in: packages/worker/src/worker.ts:598

Close the AMQP channel and connection.

Graceful shutdown in three steps: cancel every consumer (no new deliveries), drain in-flight handlers so their acks/nacks land on the still-open channel, then close the channel and release the connection.

The drain waits up to drainTimeoutMs (default DEFAULT_DRAIN_TIMEOUT_MS) for in-flight handlers. On timeout the teardown proceeds anyway — the un-acked deliveries are redelivered by the broker, preserving at-least-once semantics — and a warning is logged. Pass null to wait indefinitely.

Parameters
ParameterType
options?{ drainTimeoutMs?: number | null; }
options.drainTimeoutMs?number | null
Returns

AsyncResult<void, never>

Example
typescript
await worker.close().get();
create()
ts
static create<TContract, TCreated, TContext>(__namedParameters): AsyncResult<TypedAmqpWorker<TContract>, never>;

Defined in: packages/worker/src/worker.ts:453

Create a type-safe AMQP worker from a contract.

Connection management (including automatic reconnection) is handled internally by amqp-connection-manager via the AmqpClient. The worker will set up consumers for all contract-defined handlers asynchronously in the background once the underlying connection and channels are ready.

Connections are automatically shared across clients and workers with the same URLs and connection options, following RabbitMQ best practices.

Type Parameters
Type ParameterDefault type
TContract extends ContractDefinition-
TCreated extends Record<string, unknown> | EmptyContextEmptyContext
TContext extends Record<string, unknown> | EmptyContextTCreated
Parameters
ParameterType
__namedParametersCreateWorkerOptions<TContract, TCreated, TContext>
Returns

AsyncResult<TypedAmqpWorker<TContract>, never>

An AsyncResult that resolves to the worker. A setup/connection failure surfaces through the Defect channel (a TechnicalError cause), never a modeled Err.

Example
typescript
const result = await TypedAmqpWorker.create({
  contract: myContract,
  handlers: {
    processOrder: ({ payload }) => OkAsync(undefined),
  },
  urls: ['amqp://localhost'],
});

Type Aliases

AnyWorkerMiddleware

ts
type AnyWorkerMiddleware = WorkerMiddleware<Record<string, unknown>, Record<string, unknown>>;

Defined in: packages/worker/src/middleware.ts:102

Widened middleware shape used internally by the dispatcher, where context types have been erased. Public call sites keep the typed WorkerMiddleware form.


ConsumerOptions

ts
type ConsumerOptions = Pick<AmqpConsumeOptions, "prefetch" | "priority" | "arguments" | "consumerTag" | "exclusive">;

Defined in: packages/worker/src/worker.ts:133

Consumer options a worker handler may configure — a curated subset of the AMQP consume options. noAck and noLocal are deliberately excluded: noAck: true would silently break the worker's ack-exactly-once and retry/DLQ invariants (deliveries would be considered settled on send), and noLocal is not supported by RabbitMQ.


CreateWorkerOptions

ts
type CreateWorkerOptions<TContract, TCreated, TContext> = object;

Defined in: packages/worker/src/worker.ts:191

Options for creating a type-safe AMQP worker.

Example

typescript
const options: CreateWorkerOptions<typeof contract> = {
  contract: myContract,
  handlers: {
    // Simple handler
    processOrder: ({ payload }) => {
      console.log('Processing order:', payload.orderId);
      return OkAsync(undefined);
    },
    // Handler with prefetch configuration
    processPayment: [
      ({ payload }) => {
        console.log('Processing payment:', payload.paymentId);
        return OkAsync(undefined);
      },
      { prefetch: 10 }
    ]
  },
  urls: ['amqp://localhost'],
  defaultConsumerOptions: {
    prefetch: 5,
  },
  connectionOptions: {
    heartbeatIntervalInSeconds: 30
  },
  logger: myLogger
};

Note: Retry configuration is defined at the queue level in the contract, not at the handler level. See QueueDefinition.retry for configuration options.

Type Parameters

Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TCreated extends Record<string, unknown> | EmptyContextEmptyContext-
TContext extends TCreatedTCreated-

Properties

PropertyTypeDescriptionDefined in
connectionOptions?AmqpConnectionManagerOptionsOptional connection configuration (heartbeat, reconnect settings, etc.)packages/worker/src/worker.ts:249
connectTimeoutMs?number | nullMaximum time in ms to wait for the AMQP connection to become ready before create() resolves to a Defect (a TechnicalError cause). Defaults to 30s (the AmqpClient's DEFAULT_CONNECT_TIMEOUT_MS). Pass null to disable the timeout and let amqp-connection-manager retry indefinitely.packages/worker/src/worker.ts:269
contractTContractThe AMQP contract definition specifying consumers and their message schemaspackages/worker/src/worker.ts:197
createContext?(info) => TCreated | Promise<TCreated>Build the per-message dependency context — the seed of the middleware chain (and, without middleware, the context handlers receive directly in helpers.context). Invoked once per message after validation, so it can produce request-scoped values (correlation-id loggers, per-message transactions); close over singletons for per-worker dependencies. A rejection/throw routes the message to the DLQ as a NonRetryableError. demesne's Layer.forkScope is the recommended implementation for DI-managed graphs.packages/worker/src/worker.ts:226
defaultConsumerOptions?ConsumerOptionsOptional default consumer options applied to all consumer handlers. Handler-specific options provided in tuple form override these defaults.packages/worker/src/worker.ts:262
handlersWorkerInferHandlers<TContract, TContext>Handlers for each consumers and rpcs entry in the contract. - Regular consumers return AsyncResult<void, HandlerError>. - RPC handlers return AsyncResult<TResponse, HandlerError> where TResponse is inferred from the RPC's response message schema. When the RPC declares an errors map, the error channel additionally accepts the declared RpcError<code, data> members (otherwise it stays plain HandlerError). Handlers receive the middleware-produced context as a third argument (an empty object when no middleware is configured). Use declareHandler / declareHandlers to create handlers with full type inference.packages/worker/src/worker.ts:214
logger?LoggerOptional logger for logging message consumption and errorspackages/worker/src/worker.ts:251
maxDecompressedBytes?numberCap on the decompressed size (bytes) of a single inbound message. Guards against a decompression bomb — a few-KB payload that expands to gigabytes before schema validation runs. Over-cap messages follow the poison-message DLQ path. Defaults to DEFAULT_MAX_DECOMPRESSED_BYTES (64 MiB).packages/worker/src/worker.ts:286
middleware?| WorkerMiddleware<TCreated, TContext> | readonly AnyWorkerMiddleware[]Optional middleware wrapping every handler invocation (consumers and RPCs), applied after message validation. The chain is seeded with the createContext result (an empty object when none is configured). Accepts either a single middleware or an array (first entry = outermost, mirroring the client's interceptor arrays). The array form composes at runtime exactly like composeMiddleware(...), but cannot thread the stepwise context types across entries — when middleware accumulate typed context for the handlers, pre-compose with composeMiddleware(outermost, ..., innermost) so the chain's final context type is inferred into helpers.context. A middleware can short-circuit by returning without calling next: handler-style errors route through retry/DLQ (or a typed RPC error reply), and an Ok(value) skips the handler entirely. next({ payload }) substitutes the message payload, re-validated before the handler runs.packages/worker/src/worker.ts:245
publishTimeoutMs?number | nullMaximum time in ms a worker-side publish (retry republish, RPC reply) may sit buffered waiting for the broker before its promise settles with a timeout failure (surfaced as a Defect). Maps to amqp-connection-manager's channel-level publishTimeout. Defaults to 30s (the AmqpClient's DEFAULT_PUBLISH_TIMEOUT_MS). Pass null to disable, restoring unbounded buffering — a publish issued during an outage then never settles.packages/worker/src/worker.ts:279
telemetry?TelemetryProviderOptional telemetry provider for tracing and metrics. If not provided, uses the default provider which attempts to load OpenTelemetry. OpenTelemetry instrumentation is automatically enabled if @opentelemetry/api is installed.packages/worker/src/worker.ts:257
urlsConnectionUrl[]AMQP broker URL(s). Multiple URLs provide failover supportpackages/worker/src/worker.ts:247

EmptyContext

ts
type EmptyContext = Record<never, never>;

Defined in: packages/worker/src/middleware.ts:13

The middleware context handlers see when no middleware injects anything.

Record<never, never> rather than {} so an empty context is a real "no properties" type instead of the anything-goes empty-object type.


HandlerError

ts
type HandlerError = 
  | RetryableError
  | NonRetryableError;

Defined in: packages/worker/src/errors.ts:61

Any handler-signalled error — the union a handler may put in the Err channel of its AsyncResult. Discriminate on _tag ("@amqp-contract/RetryableError" / "@amqp-contract/NonRetryableError"), e.g. with the error matcher (matcher.with(P.tag("@amqp-contract/RetryableError"), …)).

Previously an abstract base class; now a tagged union, because unthrown's TaggedError mints a distinct base class per tag. Use isHandlerError for runtime narrowing instead of instanceof HandlerError.


Logger

ts
type Logger = object;

Defined in: packages/core/dist/index.d.mts:35

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/dist/index.d.mts:41

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/dist/index.d.mts:59

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/dist/index.d.mts:47

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/dist/index.d.mts:53

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/dist/index.d.mts:15

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/dist/index.d.mts:16

TelemetryProvider

ts
type TelemetryProvider = object;

Defined in: packages/core/dist/telemetry-BtfQAlOJ.d.mts:59

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/dist/telemetry-BtfQAlOJ.d.mts:74
getConsumeLatencyHistogram() => Histogram | undefinedGet a histogram for consume/process latency. Returns undefined if OpenTelemetry is not available.packages/core/dist/telemetry-BtfQAlOJ.d.mts: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/dist/telemetry-BtfQAlOJ.d.mts:90
getPublishCounter() => Counter | undefinedGet a counter for messages published. Returns undefined if OpenTelemetry is not available.packages/core/dist/telemetry-BtfQAlOJ.d.mts:69
getPublishLatencyHistogram() => Histogram | undefinedGet a histogram for publish latency. Returns undefined if OpenTelemetry is not available.packages/core/dist/telemetry-BtfQAlOJ.d.mts:79
getTracer() => Tracer | undefinedGet a tracer instance for creating spans. Returns undefined if OpenTelemetry is not available.packages/core/dist/telemetry-BtfQAlOJ.d.mts:64

WorkerConsumedMessage

ts
type WorkerConsumedMessage<TPayload, THeaders> = object;

Defined in: packages/worker/src/types.ts:228

A consumed message containing parsed payload and headers.

This type represents the first argument passed to consumer handlers. It contains the validated payload and (if defined in the message schema) the validated headers.

Example

typescript
const handler = declareHandler(contract, 'processOrder', (message, rawMessage) => {
  console.log(message.payload.orderId);  // Typed payload
  console.log(message.headers?.priority); // Typed headers (if defined)
  console.log(rawMessage.fields.deliveryTag); // Raw AMQP message
  return OkAsync(undefined);
});

Type Parameters

Type ParameterDefault typeDescription
TPayload-The inferred payload type from the message schema
THeadersundefinedThe inferred headers type from the message schema (undefined if not defined)

Properties

PropertyTypeDescriptionDefined in
headersTHeaders extends undefined ? undefined : THeadersThe validated message headers (present only when headers schema is defined)packages/worker/src/types.ts:232
payloadTPayloadThe validated message payloadpackages/worker/src/types.ts:230

WorkerCreateContextInfo

ts
type WorkerCreateContextInfo = object;

Defined in: packages/worker/src/worker.ts:100

Per-message information handed to the createContext factory — enough to derive request-scoped dependencies (correlation-id loggers, per-message transactions) without closing over the dispatch loop.

Properties

PropertyTypeDescriptionDefined in
handlerNamestringThe consumers / rpcs key being dispatched.packages/worker/src/worker.ts:102
isRpcbooleanTrue when the handler is an RPC server.packages/worker/src/worker.ts:104
messageobjectThe validated message (payload and headers already schema-checked).packages/worker/src/worker.ts:106
message.headersunknown-packages/worker/src/worker.ts:106
message.payloadunknown-packages/worker/src/worker.ts:106
rawMessageConsumeMessageThe raw amqplib message.packages/worker/src/worker.ts:108

WorkerHandlerHelpers

ts
type WorkerHandlerHelpers<TContext, TErrors> = object;

Defined in: packages/worker/src/types.ts:152

The helpers object every handler receives as its third argument: the middleware-produced context and (for RPC handlers with a declared errors map) the typed error constructors.

Type Parameters

Type ParameterDefault type
TContext extends Record<string, unknown> | EmptyContextEmptyContext
TErrorsEmptyContext

Properties

PropertyModifierTypeDescriptionDefined in
contextreadonlyTContextContext produced by createContext and the middleware chain.packages/worker/src/types.ts:157
errorsreadonlyTErrorsTyped constructors for the contract-declared errors (empty for consumers).packages/worker/src/types.ts:159

WorkerInferConsumedMessage

ts
type WorkerInferConsumedMessage<TContract, TName> = WorkerConsumedMessage<WorkerInferConsumerPayload<TContract, TName>, WorkerInferConsumerHeaders<TContract, TName>>;

Defined in: packages/worker/src/types.ts:238

Infer the full consumed message type for a regular consumer.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferConsumerNames<TContract>

WorkerInferConsumerHandler

ts
type WorkerInferConsumerHandler<TContract, TName, TContext> = (message, rawMessage, helpers) => AsyncResult<void, HandlerError>;

Defined in: packages/worker/src/types.ts:276

Handler signature for a regular consumer (event/command). Returns AsyncResult<void, HandlerError> — there is no response message.

Type Parameters

Type ParameterDefault type
TContract extends ContractDefinition-
TName extends InferConsumerNames<TContract>-
TContext extends Record<string, unknown> | EmptyContextEmptyContext

Parameters

ParameterType
messageWorkerInferConsumedMessage<TContract, TName>
rawMessageConsumeMessage
helpersWorkerHandlerHelpers<TContext>

Returns

AsyncResult<void, HandlerError>


WorkerInferConsumerHandlerEntry

ts
type WorkerInferConsumerHandlerEntry<TContract, TName, TContext> = 
  | WorkerInferConsumerHandler<TContract, TName, TContext>
  | readonly [WorkerInferConsumerHandler<TContract, TName, TContext>, ConsumerOptions];

Defined in: packages/worker/src/types.ts:312

Handler entry for a regular consumer — function or [handler, options].

Type Parameters

Type ParameterDefault type
TContract extends ContractDefinition-
TName extends InferConsumerNames<TContract>-
TContext extends Record<string, unknown> | EmptyContextEmptyContext

WorkerInferConsumerHeaders

ts
type WorkerInferConsumerHeaders<TContract, TName> = ConsumerInferHeadersOutput<InferConsumer<TContract, TName>>;

Defined in: packages/worker/src/types.ts:86

Infer the headers type for a regular consumer. Returns undefined if no headers schema is defined.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferConsumerNames<TContract>

WorkerInferHandlers

ts
type WorkerInferHandlers<TContract, TContext> = [InferConsumerNames<TContract>] extends [never] ? object : { [K in InferConsumerNames<TContract>]: WorkerInferConsumerHandlerEntry<TContract, K, TContext> } & [InferRpcNames<TContract>] extends [never] ? object : { [K in InferRpcNames<TContract>]: WorkerInferRpcHandlerEntry<TContract, K, TContext> };

Defined in: packages/worker/src/types.ts:351

All handlers for a contract: one entry per consumers key plus one entry per rpcs key. The two name spaces are disjoint so the resulting object type is unambiguous.

TContext is the context produced by the worker's middleware chain; the third handler argument is typed with it.

Type Parameters

Type ParameterDefault type
TContract extends ContractDefinition-
TContext extends Record<string, unknown> | EmptyContextEmptyContext

Example

typescript
const handlers: WorkerInferHandlers<typeof contract> = {
  processOrder: ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
  calculate: ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
};

WorkerInferRpcConsumedMessage

ts
type WorkerInferRpcConsumedMessage<TContract, TName> = WorkerConsumedMessage<WorkerInferRpcRequest<TContract, TName>, WorkerInferRpcHeaders<TContract, TName>>;

Defined in: packages/worker/src/types.ts:250

Infer the consumed message type for an RPC handler — payload + headers from the request side of the RPC.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerInferRpcErrorConstructors

ts
type WorkerInferRpcErrorConstructors<TContract, TName> = InferRpc<TContract, TName> extends RpcDefinition<MessageDefinition, MessageDefinition, QueueDefinition, infer TErrors> ? TErrors extends RpcErrorMap ? { [K in keyof TErrors & string]: (data: InferSchemaInput<TErrors[K]["data"]>, message?: string) => RpcError<K, InferSchemaInput<TErrors[K]["data"]>> } : EmptyContext : EmptyContext;

Defined in: packages/worker/src/types.ts:127

Typed constructors for an RPC's declared errors, handed to the handler via its helpers argument: errors.ORDER_NOT_FOUND({ orderId }) builds the RpcError with per-code data inference and autocomplete — the constructor-bag form of the free rpcError(code, data) factory (org DNA, mirroring temporal-contract's helpers.errors).

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerInferRpcErrors

ts
type WorkerInferRpcErrors<TContract, TName> = InferRpc<TContract, TName> extends RpcDefinition<MessageDefinition, MessageDefinition, QueueDefinition, infer TErrors> ? TErrors extends RpcErrorMap ? { [K in keyof TErrors & string]: RpcError<K, InferSchemaInput<TErrors[K]["data"]>> }[keyof TErrors & string] : never : never;

Defined in: packages/worker/src/types.ts:169

Infer the typed error union for an RPC handler — one RpcError<code, data> member per entry in the RPC's errors map, with data typed as the declared schema's input (the worker validates before replying). Resolves to never when the RPC declares no errors, leaving the handler's error channel as plain HandlerError.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerInferRpcHandler

ts
type WorkerInferRpcHandler<TContract, TName, TContext> = (message, rawMessage, helpers) => AsyncResult<WorkerInferRpcResponse<TContract, TName>, 
  | HandlerError
  | WorkerInferRpcErrors<TContract, TName>>;

Defined in: packages/worker/src/types.ts:296

Handler signature for an RPC. Returns AsyncResult<TResponse, HandlerError | RpcError> where TResponse is the inferred response payload and the RpcError members come from the RPC's declared errors map (absent when none are declared). The worker validates the response against the RPC's response schema and publishes it back to msg.properties.replyTo with the same correlationId; a declared RpcError is validated, published as an error reply, and the request is acked (business errors are not retried).

Type Parameters

Type ParameterDefault type
TContract extends ContractDefinition-
TName extends InferRpcNames<TContract>-
TContext extends Record<string, unknown> | EmptyContextEmptyContext

Parameters

ParameterType
messageWorkerInferRpcConsumedMessage<TContract, TName>
rawMessageConsumeMessage
helpersWorkerHandlerHelpers<TContext, WorkerInferRpcErrorConstructors<TContract, TName>>

Returns

AsyncResult<WorkerInferRpcResponse<TContract, TName>, | HandlerError | WorkerInferRpcErrors<TContract, TName>>


WorkerInferRpcHandlerEntry

ts
type WorkerInferRpcHandlerEntry<TContract, TName, TContext> = 
  | WorkerInferRpcHandler<TContract, TName, TContext>
  | readonly [WorkerInferRpcHandler<TContract, TName, TContext>, ConsumerOptions];

Defined in: packages/worker/src/types.ts:323

Handler entry for an RPC — function or [handler, options].

Type Parameters

Type ParameterDefault type
TContract extends ContractDefinition-
TName extends InferRpcNames<TContract>-
TContext extends Record<string, unknown> | EmptyContextEmptyContext

WorkerInferRpcHeaders

ts
type WorkerInferRpcHeaders<TContract, TName> = InferRpc<TContract, TName> extends RpcDefinition<infer TRequest, MessageDefinition> ? TRequest extends MessageDefinition<infer _TPayload, infer THeaders> ? THeaders extends StandardSchemaV1<Record<string, unknown>> ? InferSchemaOutput<THeaders> : undefined : undefined : undefined;

Defined in: packages/worker/src/types.ts:108

Infer the request headers type for an RPC. Returns undefined unless the RPC's request MessageDefinition declares a headers schema.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerInferRpcRequest

ts
type WorkerInferRpcRequest<TContract, TName> = InferRpc<TContract, TName> extends RpcDefinition<infer TRequest, MessageDefinition> ? TRequest extends MessageDefinition ? InferSchemaOutput<TRequest["payload"]> : never : never;

Defined in: packages/worker/src/types.ts:94

Infer the request payload type for an RPC.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerInferRpcResponse

ts
type WorkerInferRpcResponse<TContract, TName> = InferRpc<TContract, TName> extends RpcDefinition<MessageDefinition, infer TResponse> ? TResponse extends MessageDefinition ? InferSchemaInput<TResponse["payload"]> : never : never;

Defined in: packages/worker/src/types.ts:195

Infer the response payload type for an RPC. The handler must return an AsyncResult<TResponse, HandlerError> matching this shape.

Typed as the schema's input — the handler supplies the pre-validation shape (defaults optional, transforms not yet applied); the worker validates against the response schema before publishing the reply. Same convention as RPC error data.

Type Parameters

Type Parameter
TContract extends ContractDefinition
TName extends InferRpcNames<TContract>

WorkerMiddleware

ts
type WorkerMiddleware<TContextIn, TContextOut> = (args, next) => AsyncResult<unknown, HandlerError | RpcError>;

Defined in: packages/worker/src/middleware.ts:89

A worker middleware: wraps every handler invocation (consumers and RPCs) after message validation.

Middleware follow the guard-and-narrow pattern: check something, then call next({ context }) to run the rest of the chain with typed fields injected into the handler's third argument — or short-circuit by returning without calling next:

  • Err(retryable(...)) / Err(nonRetryable(...)) routes through the normal retry/DLQ pipeline, exactly as if the handler had returned it.
  • Err(rpcError(code, data)) (on an RPC with a declared errors map) publishes a typed error reply.
  • Ok(value) skips the handler entirely; for an RPC, value is validated against the response schema and published as the reply (cache pattern).

Type Parameters

Type ParameterDefault typeDescription
TContextIn extends Record<string, unknown> | EmptyContextEmptyContextcontext provided by outer middleware
TContextOut extends TContextInTContextIncontext this middleware passes downstream

Parameters

ParameterType
argsWorkerMiddlewareArgs<TContextIn>
nextWorkerMiddlewareNext<TContextOut>

Returns

AsyncResult<unknown, HandlerError | RpcError>

Example

typescript
import { declareMiddleware, nonRetryable } from '@amqp-contract/worker';
import { ErrAsync } from 'unthrown';

const auth = declareMiddleware<EmptyContext, { tenantId: string }>((args, next) => {
  const tenantId = args.rawMessage.properties.headers?.['x-tenant-id'];
  if (typeof tenantId !== 'string') {
    return ErrAsync(nonRetryable('Missing x-tenant-id header'));
  }
  return next({ context: { tenantId } });
});

WorkerMiddlewareArgs

ts
type WorkerMiddlewareArgs<TContextIn> = object;

Defined in: packages/worker/src/middleware.ts:20

Arguments passed to a worker middleware — everything the downstream handler will receive, plus dispatch metadata and the context accumulated by outer middleware.

Type Parameters

Type Parameter
TContextIn extends Record<string, unknown> | EmptyContext

Properties

PropertyTypeDescriptionDefined in
contextTContextInContext accumulated by outer middleware (empty for the outermost one).packages/worker/src/middleware.ts:30
handlerNamestringThe consumers / rpcs key being dispatched.packages/worker/src/middleware.ts:26
isRpcbooleanTrue when the handler is an RPC server (its result is published as a reply).packages/worker/src/middleware.ts:28
messageobjectThe validated message (payload and headers already schema-checked).packages/worker/src/middleware.ts:22
message.headersunknown-packages/worker/src/middleware.ts:22
message.payloadunknown-packages/worker/src/middleware.ts:22
rawMessageConsumeMessageThe raw amqplib message (delivery tag, AMQP properties, raw headers, …).packages/worker/src/middleware.ts:24

WorkerMiddlewareNext

ts
type WorkerMiddlewareNext<TContextOut> = (opts?) => AsyncResult<unknown, HandlerError | RpcError>;

Defined in: packages/worker/src/middleware.ts:50

Continuation a middleware calls to run the rest of the chain (inner middleware, then the handler). opts.context is the full context the downstream sees — pass { ...args.context, ...injected } (or just the injected fields; the dispatcher merges over the current context either way).

opts.payload substitutes the message payload seen downstream. Inner middleware observe the substituted payload as-is; the dispatcher re-validates it against the consumer's payload schema before the handler runs — an invalid substitution fails terminally as a NonRetryableError (DLQ), so middleware cannot smuggle unvalidated data past the contract boundary.

The returned AsyncResult carries the handler outcome: undefined for a regular consumer, the (not yet validated) response for an RPC. A middleware can transform or inspect it before returning.

Type Parameters

Type Parameter
TContextOut extends Record<string, unknown> | EmptyContext

Parameters

ParameterType
opts?{ context?: TContextOut; payload?: unknown; }
opts.context?TContextOut
opts.payload?unknown

Returns

AsyncResult<unknown, HandlerError | RpcError>

Variables

DEFAULT_DRAIN_TIMEOUT_MS

ts
const DEFAULT_DRAIN_TIMEOUT_MS: 30000 = 30_000;

Defined in: packages/worker/src/worker.ts:144

Default time close() waits for in-flight handlers before tearing the channel down anyway. Finite by default so a hung handler cannot wedge shutdown — the un-acked deliveries are redelivered by the broker (at-least-once semantics). Pass drainTimeoutMs: null to wait forever.

Functions

composeMiddleware()

Call Signature

ts
function composeMiddleware<TIn, TA>(m1): WorkerMiddleware<TIn, TA>;

Defined in: packages/worker/src/middleware.ts:140

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
Returns

WorkerMiddleware<TIn, TA>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB>(m1, m2): WorkerMiddleware<TIn, TB>;

Defined in: packages/worker/src/middleware.ts:144

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
Returns

WorkerMiddleware<TIn, TB>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC>(
   m1, 
   m2, 
   m3): WorkerMiddleware<TIn, TC>;

Defined in: packages/worker/src/middleware.ts:149

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
Returns

WorkerMiddleware<TIn, TC>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC, TD>(
   m1, 
   m2, 
   m3, 
   m4): WorkerMiddleware<TIn, TD>;

Defined in: packages/worker/src/middleware.ts:159

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
TD extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
m4WorkerMiddleware<TC, TD>
Returns

WorkerMiddleware<TIn, TD>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC, TD, TE>(
   m1, 
   m2, 
   m3, 
   m4, 
   m5): WorkerMiddleware<TIn, TE>;

Defined in: packages/worker/src/middleware.ts:171

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
TD extends Record<string, unknown> | EmptyContext
TE extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
m4WorkerMiddleware<TC, TD>
m5WorkerMiddleware<TD, TE>
Returns

WorkerMiddleware<TIn, TE>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC, TD, TE, TF>(
   m1, 
   m2, 
   m3, 
   m4, 
   m5, 
   m6): WorkerMiddleware<TIn, TF>;

Defined in: packages/worker/src/middleware.ts:185

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
TD extends Record<string, unknown> | EmptyContext
TE extends Record<string, unknown> | EmptyContext
TF extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
m4WorkerMiddleware<TC, TD>
m5WorkerMiddleware<TD, TE>
m6WorkerMiddleware<TE, TF>
Returns

WorkerMiddleware<TIn, TF>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC, TD, TE, TF, TG>(
   m1, 
   m2, 
   m3, 
   m4, 
   m5, 
   m6, 
   m7): WorkerMiddleware<TIn, TG>;

Defined in: packages/worker/src/middleware.ts:201

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
TD extends Record<string, unknown> | EmptyContext
TE extends Record<string, unknown> | EmptyContext
TF extends Record<string, unknown> | EmptyContext
TG extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
m4WorkerMiddleware<TC, TD>
m5WorkerMiddleware<TD, TE>
m6WorkerMiddleware<TE, TF>
m7WorkerMiddleware<TF, TG>
Returns

WorkerMiddleware<TIn, TG>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

Call Signature

ts
function composeMiddleware<TIn, TA, TB, TC, TD, TE, TF, TG, TH>(
   m1, 
   m2, 
   m3, 
   m4, 
   m5, 
   m6, 
   m7, 
   m8): WorkerMiddleware<TIn, TH>;

Defined in: packages/worker/src/middleware.ts:219

Compose middleware left-to-right: the first argument is the outermost (runs first, sees the emptiest context), the last is the innermost (its injected context is what handlers receive). Context types accumulate across the chain — each middleware's TContextIn must match what the previous one produced.

Typed overloads cover up to 8 middleware. For longer chains, nest: a composed chain is itself a WorkerMiddleware<EmptyContext, T> and can be the first argument of an outer composeMiddleware call — composeMiddleware(composeMiddleware(a, ..., h), i, j) — preserving context-type accumulation at any length.

Type Parameters
Type Parameter
TIn extends Record<string, unknown> | EmptyContext
TA extends Record<string, unknown> | EmptyContext
TB extends Record<string, unknown> | EmptyContext
TC extends Record<string, unknown> | EmptyContext
TD extends Record<string, unknown> | EmptyContext
TE extends Record<string, unknown> | EmptyContext
TF extends Record<string, unknown> | EmptyContext
TG extends Record<string, unknown> | EmptyContext
TH extends Record<string, unknown> | EmptyContext
Parameters
ParameterType
m1WorkerMiddleware<TIn, TA>
m2WorkerMiddleware<TA, TB>
m3WorkerMiddleware<TB, TC>
m4WorkerMiddleware<TC, TD>
m5WorkerMiddleware<TD, TE>
m6WorkerMiddleware<TE, TF>
m7WorkerMiddleware<TF, TG>
m8WorkerMiddleware<TG, TH>
Returns

WorkerMiddleware<TIn, TH>

Example
typescript
const middleware = composeMiddleware(logging, auth, idempotency);
// handlers receive the context injected by all three

declareHandler()

Call Signature

ts
function declareHandler<TContract, TName, TContext>(
   contract, 
   name, 
   handler): WorkerInferConsumerHandlerEntry<TContract, TName, TContext>;

Defined in: packages/worker/src/handlers.ts:206

Define a type-safe handler for a specific consumer or RPC in a contract.

Recommended: This function creates handlers that return AsyncResult<void, HandlerError> (consumers) or AsyncResult<TResponse, HandlerError> (RPCs), providing explicit error handling and better control over retry behavior.

Supports two patterns:

  1. Simple handler: just the function
  2. Handler with options: [handler, { prefetch: 10 }]
Type Parameters
Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TName extends string | number | symbol-The consumer or RPC name from the contract
TContext extends Record<string, unknown> | EmptyContextEmptyContext-
Parameters
ParameterTypeDescription
contractTContractThe contract definition containing the consumer or RPC
nameTNameThe name of the consumer or RPC from the contract
handlerWorkerInferConsumerHandler<TContract, TName, TContext>The handler function — for consumers, returns AsyncResult<void, HandlerError>; for RPCs, returns AsyncResult<TResponse, HandlerError>.
Returns

WorkerInferConsumerHandlerEntry<TContract, TName, TContext>

A type-safe handler that can be used with TypedAmqpWorker

Examples

Consumer handler

typescript
import { declareHandler, RetryableError, NonRetryableError } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const processOrderHandler = declareHandler(
  orderContract,
  'processOrder',
  ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
);

RPC handler

typescript
const calculateHandler = declareHandler(
  rpcContract,
  'calculate',
  ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
);

Call Signature

ts
function declareHandler<TContract, TName, TContext>(
   contract, 
   name, 
   handler, 
   options): WorkerInferConsumerHandlerEntry<TContract, TName, TContext>;

Defined in: packages/worker/src/handlers.ts:215

Define a type-safe handler for a specific consumer or RPC in a contract.

Recommended: This function creates handlers that return AsyncResult<void, HandlerError> (consumers) or AsyncResult<TResponse, HandlerError> (RPCs), providing explicit error handling and better control over retry behavior.

Supports two patterns:

  1. Simple handler: just the function
  2. Handler with options: [handler, { prefetch: 10 }]
Type Parameters
Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TName extends string | number | symbol-The consumer or RPC name from the contract
TContext extends Record<string, unknown> | EmptyContextEmptyContext-
Parameters
ParameterTypeDescription
contractTContractThe contract definition containing the consumer or RPC
nameTNameThe name of the consumer or RPC from the contract
handlerWorkerInferConsumerHandler<TContract, TName, TContext>The handler function — for consumers, returns AsyncResult<void, HandlerError>; for RPCs, returns AsyncResult<TResponse, HandlerError>.
optionsConsumerOptionsOptional consumer options (prefetch)
Returns

WorkerInferConsumerHandlerEntry<TContract, TName, TContext>

A type-safe handler that can be used with TypedAmqpWorker

Examples

Consumer handler

typescript
import { declareHandler, RetryableError, NonRetryableError } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const processOrderHandler = declareHandler(
  orderContract,
  'processOrder',
  ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
);

RPC handler

typescript
const calculateHandler = declareHandler(
  rpcContract,
  'calculate',
  ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
);

Call Signature

ts
function declareHandler<TContract, TName, TContext>(
   contract, 
   name, 
   handler): WorkerInferRpcHandlerEntry<TContract, TName, TContext>;

Defined in: packages/worker/src/handlers.ts:225

Define a type-safe handler for a specific consumer or RPC in a contract.

Recommended: This function creates handlers that return AsyncResult<void, HandlerError> (consumers) or AsyncResult<TResponse, HandlerError> (RPCs), providing explicit error handling and better control over retry behavior.

Supports two patterns:

  1. Simple handler: just the function
  2. Handler with options: [handler, { prefetch: 10 }]
Type Parameters
Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TName extends string | number | symbol-The consumer or RPC name from the contract
TContext extends Record<string, unknown> | EmptyContextEmptyContext-
Parameters
ParameterTypeDescription
contractTContractThe contract definition containing the consumer or RPC
nameTNameThe name of the consumer or RPC from the contract
handlerWorkerInferRpcHandler<TContract, TName, TContext>The handler function — for consumers, returns AsyncResult<void, HandlerError>; for RPCs, returns AsyncResult<TResponse, HandlerError>.
Returns

WorkerInferRpcHandlerEntry<TContract, TName, TContext>

A type-safe handler that can be used with TypedAmqpWorker

Examples

Consumer handler

typescript
import { declareHandler, RetryableError, NonRetryableError } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const processOrderHandler = declareHandler(
  orderContract,
  'processOrder',
  ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
);

RPC handler

typescript
const calculateHandler = declareHandler(
  rpcContract,
  'calculate',
  ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
);

Call Signature

ts
function declareHandler<TContract, TName, TContext>(
   contract, 
   name, 
   handler, 
   options): WorkerInferRpcHandlerEntry<TContract, TName, TContext>;

Defined in: packages/worker/src/handlers.ts:234

Define a type-safe handler for a specific consumer or RPC in a contract.

Recommended: This function creates handlers that return AsyncResult<void, HandlerError> (consumers) or AsyncResult<TResponse, HandlerError> (RPCs), providing explicit error handling and better control over retry behavior.

Supports two patterns:

  1. Simple handler: just the function
  2. Handler with options: [handler, { prefetch: 10 }]
Type Parameters
Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TName extends string | number | symbol-The consumer or RPC name from the contract
TContext extends Record<string, unknown> | EmptyContextEmptyContext-
Parameters
ParameterTypeDescription
contractTContractThe contract definition containing the consumer or RPC
nameTNameThe name of the consumer or RPC from the contract
handlerWorkerInferRpcHandler<TContract, TName, TContext>The handler function — for consumers, returns AsyncResult<void, HandlerError>; for RPCs, returns AsyncResult<TResponse, HandlerError>.
optionsConsumerOptionsOptional consumer options (prefetch)
Returns

WorkerInferRpcHandlerEntry<TContract, TName, TContext>

A type-safe handler that can be used with TypedAmqpWorker

Examples

Consumer handler

typescript
import { declareHandler, RetryableError, NonRetryableError } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const processOrderHandler = declareHandler(
  orderContract,
  'processOrder',
  ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
);

RPC handler

typescript
const calculateHandler = declareHandler(
  rpcContract,
  'calculate',
  ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
);

declareHandlers()

ts
function declareHandlers<TContract, TContext>(contract, handlers): WorkerInferHandlers<TContract, TContext>;

Defined in: packages/worker/src/handlers.ts:287

Define multiple type-safe handlers for consumers and RPCs in a contract.

Recommended: This function creates handlers that return AsyncResult<void, HandlerError> (consumers) or AsyncResult<TResponse, HandlerError> (RPCs), providing explicit error handling and better control over retry behavior.

The handlers object must contain exactly one entry per consumers and rpcs key in the contract — see WorkerInferHandlers.

Type Parameters

Type ParameterDefault typeDescription
TContract extends ContractDefinition-The contract definition type
TContext extends Record<string, unknown> | EmptyContextEmptyContext-

Parameters

ParameterTypeDescription
contractTContractThe contract definition containing the consumers and RPCs
handlersWorkerInferHandlers<TContract, TContext>An object with handler functions for each consumer and RPC

Returns

WorkerInferHandlers<TContract, TContext>

A type-safe handlers object that can be used with TypedAmqpWorker

Example

typescript
import { declareHandlers, RetryableError } from '@amqp-contract/worker';
import { fromPromise, OkAsync } from 'unthrown';

const handlers = declareHandlers(orderContract, {
  processOrder: ({ payload }) =>
    fromPromise(
      processPayment(payload),
      (error) => new RetryableError('Payment failed', error),
    ).map(() => undefined),
  calculate: ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
});

declareMiddleware()

ts
function declareMiddleware<TContextIn, TContextOut>(middleware): WorkerMiddleware<TContextIn, TContextOut>;

Defined in: packages/worker/src/middleware.ts:112

Identity helper that pins a middleware's context types without a variable annotation — declareMiddleware<In, Out>(fn) reads better than const mw: WorkerMiddleware<In, Out> = fn.

Type Parameters

Type ParameterDefault type
TContextIn extends Record<string, unknown> | EmptyContextEmptyContext
TContextOut extends Record<string, unknown> | EmptyContextTContextIn

Parameters

ParameterType
middlewareWorkerMiddleware<TContextIn, TContextOut>

Returns

WorkerMiddleware<TContextIn, TContextOut>


isHandlerError()

ts
function isHandlerError(error): error is HandlerError;

Defined in: packages/worker/src/errors.ts:137

Type guard to check if an error is any HandlerError (RetryableError or NonRetryableError).

Parameters

ParameterTypeDescription
errorunknownThe error to check

Returns

error is HandlerError

True if the error is a HandlerError

Example

typescript
import { isHandlerError } from '@amqp-contract/worker';

function handleError(error: unknown) {
  if (isHandlerError(error)) {
    // error is RetryableError | NonRetryableError
    console.log('Handler error:', error.name, error.message);
  }
}

isMessageValidationError()

ts
function isMessageValidationError(error): error is MessageValidationError;

Defined in: packages/core/dist/index.d.mts:468

Type guard to check if an error is a MessageValidationError.

Parameters

ParameterType
errorunknown

Returns

error is MessageValidationError


isNonRetryableError()

ts
function isNonRetryableError(error): error is NonRetryableError;

Defined in: packages/worker/src/errors.ts:115

Type guard to check if an error is a NonRetryableError.

Use this to check error types in catch blocks or error handlers.

Parameters

ParameterTypeDescription
errorunknownThe error to check

Returns

error is NonRetryableError

True if the error is a NonRetryableError

Example

typescript
import { isNonRetryableError } from '@amqp-contract/worker';

try {
  await processMessage();
} catch (error) {
  if (isNonRetryableError(error)) {
    console.log('Will not retry:', error.message);
  }
}

isRetryableError()

ts
function isRetryableError(error): error is RetryableError;

Defined in: packages/worker/src/errors.ts:90

Type guard to check if an error is a RetryableError.

Use this to check error types in catch blocks or error handlers.

Parameters

ParameterTypeDescription
errorunknownThe error to check

Returns

error is RetryableError

True if the error is a RetryableError

Example

typescript
import { isRetryableError } from '@amqp-contract/worker';

try {
  await processMessage();
} catch (error) {
  if (isRetryableError(error)) {
    console.log('Will retry:', error.message);
  } else {
    console.log('Permanent failure:', error);
  }
}

isRpcError()

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

Defined in: packages/core/dist/index.d.mts:513

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/dist/index.d.mts:464

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


nonRetryable()

ts
function nonRetryable(message, cause?): NonRetryableError;

Defined in: packages/worker/src/errors.ts:200

Create a NonRetryableError with less verbosity.

This is a shorthand factory function for creating NonRetryableError instances. Use it for cleaner error creation in handlers.

Parameters

ParameterTypeDescription
messagestringError message describing the failure
cause?unknownOptional underlying error that caused this failure

Returns

NonRetryableError

A new NonRetryableError instance

Example

typescript
import { nonRetryable } from '@amqp-contract/worker';
import { ErrAsync, OkAsync } from 'unthrown';

const handler = ({ payload }) => {
  if (!isValidPayload(payload)) {
    return ErrAsync(nonRetryable('Invalid payload format'));
  }
  return OkAsync(undefined);
};

// Equivalent to:
// return ErrAsync(new NonRetryableError('Invalid payload format'));

qualifyNonRetryable()

ts
function qualifyNonRetryable(message): (cause) => NonRetryableError;

Defined in: packages/worker/src/errors.ts:248

Build a fromPromise qualifier that wraps the rejection cause in a NonRetryableError with the given message — the permanent-failure counterpart of qualifyRetryable.

Parameters

ParameterType
messagestring

Returns

(cause) => NonRetryableError

Example

typescript
import { qualifyNonRetryable } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const handler = ({ payload }) =>
  fromPromise(chargeCard(payload), qualifyNonRetryable('Card permanently declined'))
    .map(() => undefined);

qualifyRetryable()

ts
function qualifyRetryable(message): (cause) => RetryableError;

Defined in: packages/worker/src/errors.ts:229

Build a fromPromise qualifier that wraps the rejection cause in a RetryableError with the given message.

fromPromise requires a qualify mapper as its second argument; writing it by hand every time ((e) => retryable("...", e)) is the most re-introduced mistake in this codebase. The factory removes the boilerplate:

Parameters

ParameterType
messagestring

Returns

(cause) => RetryableError

Example

typescript
import { qualifyRetryable } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const handler = ({ payload }) =>
  fromPromise(processPayment(payload), qualifyRetryable('Payment service unavailable'))
    .map(() => undefined);

// Equivalent to:
// fromPromise(processPayment(payload), (e) => retryable('Payment service unavailable', e))

retryable()

ts
function retryable(message, cause?): RetryableError;

Defined in: packages/worker/src/errors.ts:170

Create a RetryableError with less verbosity.

This is a shorthand factory function for creating RetryableError instances. Use it for cleaner error creation in handlers.

Parameters

ParameterTypeDescription
messagestringError message describing the failure
cause?unknownOptional underlying error that caused this failure

Returns

RetryableError

A new RetryableError instance

Example

typescript
import { retryable } from '@amqp-contract/worker';
import { fromPromise } from 'unthrown';

const handler = ({ payload }) =>
  fromPromise(
    processPayment(payload),
    (e) => retryable('Payment service unavailable', e),
  ).map(() => undefined);

// Equivalent to:
// fromPromise(processPayment(payload), (e) => new RetryableError('...', e))

rpcError()

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

Defined in: packages/core/dist/index.d.mts:539

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 }));
  }
  // ...
};

Released under the MIT License.