exportasyncfunctionassertMessageProducedAndConsumed(container:StartedKafkaContainer,additionalKafkaConfig:Partial<KafkaJS.KafkaConfig>={},additionalGlobalConfig:Partial<GlobalConfig>={}){// Static loading crashes Bun before BUN_CI can skip the dependent tests.const{KafkaJS:KafkaJSClient}=awaitimport("@confluentinc/kafka-javascript");constbrokers=[`${container.getHost()}:${container.getMappedPort(9093)}`];constkafka=newKafkaJSClient.Kafka({kafkaJS:{logLevel:KafkaJSClient.logLevel.ERROR,brokers,...additionalKafkaConfig,},...additionalGlobalConfig,});constproducer=kafka.producer();awaitproducer.connect();constconsumer=kafka.consumer({kafkaJS:{groupId:"test-group",fromBeginning:true}});awaitconsumer.connect();awaitproducer.send({topic:"test-topic",messages:[{value:"test message"}]});awaitconsumer.subscribe({topic:"test-topic"});constconsumedMessage=awaitnewPromise((resolve)=>consumer.run({eachMessage:async({message})=>resolve(message.value?.toString()),}));expect(consumedMessage).toBe("test message");awaitconsumer.disconnect();awaitproducer.disconnect();}
exportasyncfunctionassertMessageProducedAndConsumed(container:StartedKafkaContainer,additionalKafkaConfig:Partial<KafkaJS.KafkaConfig>={},additionalGlobalConfig:Partial<GlobalConfig>={}){// Static loading crashes Bun before BUN_CI can skip the dependent tests.const{KafkaJS:KafkaJSClient}=awaitimport("@confluentinc/kafka-javascript");constbrokers=[`${container.getHost()}:${container.getMappedPort(9093)}`];constkafka=newKafkaJSClient.Kafka({kafkaJS:{logLevel:KafkaJSClient.logLevel.ERROR,brokers,...additionalKafkaConfig,},...additionalGlobalConfig,});constproducer=kafka.producer();awaitproducer.connect();constconsumer=kafka.consumer({kafkaJS:{groupId:"test-group",fromBeginning:true}});awaitconsumer.connect();awaitproducer.send({topic:"test-topic",messages:[{value:"test message"}]});awaitconsumer.subscribe({topic:"test-topic"});constconsumedMessage=awaitnewPromise((resolve)=>consumer.run({eachMessage:async({message})=>resolve(message.value?.toString()),}));expect(consumedMessage).toBe("test message");awaitconsumer.disconnect();awaitproducer.disconnect();}
exportasyncfunctionassertMessageProducedAndConsumed(container:StartedKafkaContainer,additionalKafkaConfig:Partial<KafkaJS.KafkaConfig>={},additionalGlobalConfig:Partial<GlobalConfig>={}){// Static loading crashes Bun before BUN_CI can skip the dependent tests.const{KafkaJS:KafkaJSClient}=awaitimport("@confluentinc/kafka-javascript");constbrokers=[`${container.getHost()}:${container.getMappedPort(9093)}`];constkafka=newKafkaJSClient.Kafka({kafkaJS:{logLevel:KafkaJSClient.logLevel.ERROR,brokers,...additionalKafkaConfig,},...additionalGlobalConfig,});constproducer=kafka.producer();awaitproducer.connect();constconsumer=kafka.consumer({kafkaJS:{groupId:"test-group",fromBeginning:true}});awaitconsumer.connect();awaitproducer.send({topic:"test-topic",messages:[{value:"test message"}]});awaitconsumer.subscribe({topic:"test-topic"});constconsumedMessage=awaitnewPromise((resolve)=>consumer.run({eachMessage:async({message})=>resolve(message.value?.toString()),}));expect(consumedMessage).toBe("test message");awaitconsumer.disconnect();awaitproducer.disconnect();}
exportasyncfunctionassertMessageProducedAndConsumed(container:StartedKafkaContainer,additionalKafkaConfig:Partial<KafkaJS.KafkaConfig>={},additionalGlobalConfig:Partial<GlobalConfig>={}){// Static loading crashes Bun before BUN_CI can skip the dependent tests.const{KafkaJS:KafkaJSClient}=awaitimport("@confluentinc/kafka-javascript");constbrokers=[`${container.getHost()}:${container.getMappedPort(9093)}`];constkafka=newKafkaJSClient.Kafka({kafkaJS:{logLevel:KafkaJSClient.logLevel.ERROR,brokers,...additionalKafkaConfig,},...additionalGlobalConfig,});constproducer=kafka.producer();awaitproducer.connect();constconsumer=kafka.consumer({kafkaJS:{groupId:"test-group",fromBeginning:true}});awaitconsumer.connect();awaitproducer.send({topic:"test-topic",messages:[{value:"test message"}]});awaitconsumer.subscribe({topic:"test-topic"});constconsumedMessage=awaitnewPromise((resolve)=>consumer.run({eachMessage:async({message})=>resolve(message.value?.toString()),}));expect(consumedMessage).toBe("test message");awaitconsumer.disconnect();awaitproducer.disconnect();}