///
///
import amqp from 'amqplib';
import EventEmitter from 'events';
export type RabbitMQ = {
initialize: (config: RabbitMQConfig) => RabbitMQConfig;
getConfig: () => RabbitMQConfig | null;
};
export type RabbitMQConfig = {
url: string;
retryDelay: number;
maxRetries: number;
/**
* Delay (ms) before the FIRST reconnect attempt after a disconnect.
* Use `0` for an instant first attempt. Subsequent attempts grow
* exponentially from this value (see `reconnectMultiplier`).
*/
reconnectDelay: number;
/**
* Hard cap on any scheduled reconnect delay (ms). Default: 30_000.
*/
maxReconnectDelay?: number;
/**
* Growth multiplier applied per failed reconnect attempt. Default: 2.
* Set to 1 to preserve the legacy fixed-delay behaviour.
*/
reconnectMultiplier?: number;
/**
* Symmetric jitter ratio in [0, 1] applied to every scheduled delay.
* Default: 0.2 (±20%).
*/
reconnectJitterRatio?: number;
/**
* Emit an error log whenever a consumer is lost (channel closed, error,
* cancelled, or connection blocked). Off by default.
*/
healthLogging?: boolean;
};
export type GetConfigs = () => RabbitMQConfig | null;
export type QueuePublisher = {
initialize: () => Promise;
pushMessage: (queueName: string, message: Buffer) => Promise<{
ack: boolean;
err?: any;
}>;
cleanUp: (isClosureIntentional: boolean) => Promise;
isConnected: () => boolean;
};
export type QueueConsumer = {
initialize: () => Promise;
startConsumer: (queueName: string, messageProcessorFunc: (msgPayload: string, metaData?: MessageProcessorMetadata) => Promise) => Promise;
cleanUp: (isClosureIntentional: boolean) => Promise;
isConnected: () => boolean;
isHealthy: () => boolean;
getMissingQueues: () => string[];
enableLogging: () => void;
disableLogging: () => void;
};
export type BaseConnection = {
getConnection: () => amqp.Connection | null;
getChannel: () => amqp.ConfirmChannel | null;
initialize: () => Promise;
cleanUp: (isClosureIntentional?: boolean) => Promise;
event: EventEmitter;
isConnectionLive: () => boolean;
};
export type MessageProcessorMetadata = {
retryCount: number;
queueName: string;
retryDelay: number;
maxRetries: number;
reconnectDelay: number;
};