import { Kafka } from "kafkajs"; import { env } from "./env.js"; import { logger } from "../shared/observability/logger.js"; export const kafka = new Kafka({ brokers: env.KAFKA_BROKERS.split(","), clientId: "core-edu-service", }); export const producer = kafka.producer({ idempotent: true, transactionalId: "core-edu-tx", }); export const consumer = kafka.consumer({ groupId: "core-edu-group" }); let producerConnected = false; let consumerConnected = false; producer.on("producer.connect", () => { producerConnected = true; }); producer.on("producer.disconnect", () => { producerConnected = false; }); consumer.on("consumer.connect", () => { consumerConnected = true; }); consumer.on("consumer.disconnect", () => { consumerConnected = false; }); export function isKafkaConnected(): boolean { return producerConnected && consumerConnected; } export async function connectKafka(): Promise { try { await producer.connect(); await consumer.connect(); logger.info("Kafka connected"); } catch (err) { logger.warn( { err: err instanceof Error ? err.message : String(err) }, "Kafka connect failed, running without Kafka", ); } } export async function disconnectKafka(): Promise { await producer.disconnect(); await consumer.disconnect(); }