interface ProducerInterface { connect: () => Promise; disconnect: () => any; send:(arg: KafkaProducerMessageInterface) => any; } interface ConsumerInterface { connect: () => Promise; disconnect: () => any; subscribe: (arg: {topic:string|RegExp, fromBeginning: boolean}) => Promise; run: (arg: ({ eachMessage: (msg: KafkaConsumerMessageInterface) => void, })) => any; stop: () => Promise; pause: () => Promise; resume: () => Promise; } interface AdminInterface { connect: () => Promise; disconnect: () => Promise; listTopics: () => Promise; createTopics: (arg: {validateOnly?:boolean, waitForLeaders?:boolean, timeout?:number, topics:{topic:string, numPartitions?:number}[]}) => Promise; deleteTopics: (arg: {topics: string[]}) => any; } interface KafkaInterface { producer: () => ProducerInterface; consumer: ({groupId:string}) => ConsumerInterface; admin: () => AdminInterface; } //timeout: number, //used for added delay per retry interface KafkaProducerMessageInterface { topic: string, offset?: number, partition?:number, messages: { key?: string, value: string, headers?: { cascadeMetadata?: string, } }[] } interface KafkaConsumerMessageInterface { topic: string, partition: number, offset: number, message: { key?: string, value: string, headers?: { cascadeMetadata?: string, } } } interface ProducerRoute { status: string, retryLevels: number, timeoutLimit: number[], batchLimit: number[], levels: KafkaProducerMessageInterface[], topics: string[], } interface CascadeMetadata { status: string, retries: number, topicArr: string[], } type ServiceCallback = (msg: KafkaConsumerMessageInterface, resolve: RouteCallback, reject: RouteCallback) => void; type RouteCallback = (msg: KafkaConsumerMessageInterface, status?:string) => void; export { ProducerInterface, ConsumerInterface, AdminInterface, KafkaInterface, KafkaProducerMessageInterface, KafkaConsumerMessageInterface, ProducerRoute, CascadeMetadata, ServiceCallback, RouteCallback, };