///
import { Transform } from 'stream';
/**
* @callback MessageConsumerCallback
* @param message {Object}
* @param message.topic {String}
* @param message.offset {Number}
* @param message.values {String|Buffer}
* @param message.partition {Number}
* @returns {Promise|*} This function need to return a promise if the message is consumed asynchronously,
* value of any other types indicates that the message have already been consumed.
*/
/**
* @callback FailedMessageConsumerCallback
* @param error
* @param message {Object}
* @param message.topic {String}
* @param message.offset {Number}
* @param message.values {String|Buffer}
* @param message.partition {Number}
* @returns {Promise|*} This function need to return a promise if the message is consumed asynchronously,
* value of any other types indicates that the message have already been consumed.
*
*/
declare namespace ConsumeTransformStream {
interface Option {
consumeConcurrency: number;
consumeTimeout: number;
groupId: string;
messageConsumer(message: any): Promise | unknown;
failedMessageConsumer?(error: any, message: any): Promise | unknown;
}
}
/**
* @private
*/
declare class ConsumeTransformStream extends Transform {
private _options;
private _currentConsumeConcurrency;
private _concurrentPromise;
private _waitingQueue;
private _lastMessageQueuedPromise;
private _isDestroyed;
private _unhandledException?;
/**
*
* @param options.messageConsumer {MessageConsumerCallback}
* @param options {Object}
* @param options.consumeConcurrency {Number}
* @param options.consumeTimeout {Number}
* @param options.groupId {String}
* @param [options.failedMessageConsumer] {FailedMessageConsumerCallback}
*/
constructor(options: ConsumeTransformStream.Option);
_transform(message: any, encoding: any, callback: any): void;
_flush(callback: any): void;
_consumeMessage(message: any): any;
_concurrentConsumeMessage(message: any): void;
_enqueue(message: any): any;
_dequeue(): void;
_internalDestroy(e: Error): void;
}
export default ConsumeTransformStream;