From 3bb716ab8593a7a55cf72b69145fc13c694fd609 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Est=C3=AAv=C3=A3o=20Samuel=20Proc=C3=B3pio=20Amaral?= Date: Tue, 11 Aug 2026 00:15:49 -0300 Subject: [PATCH] feat(rotor): make Kafka client logging and librdkafka debug configurable --- services/rotor/__tests__/kafka-config.test.ts | 43 +++++++++++++++++ services/rotor/src/lib/kafka-config.ts | 46 +++++++++++++++---- services/rotor/src/serverEnv.ts | 2 + 3 files changed, 81 insertions(+), 10 deletions(-) create mode 100644 services/rotor/__tests__/kafka-config.test.ts diff --git a/services/rotor/__tests__/kafka-config.test.ts b/services/rotor/__tests__/kafka-config.test.ts new file mode 100644 index 000000000..ed191b520 --- /dev/null +++ b/services/rotor/__tests__/kafka-config.test.ts @@ -0,0 +1,43 @@ +import { expect, test } from "vitest"; +import { KafkaJS } from "@confluentinc/kafka-javascript"; +import { buildKafkaConfig, kafkaClientLogLevel } from "../src/lib/kafka-config"; + +const opts = { defaultAppId: "test", brokers: ["localhost:9092"] }; + +test("client log level falls back to error for missing or unknown values", () => { + expect(kafkaClientLogLevel(undefined)).toBe(KafkaJS.logLevel.ERROR); + expect(kafkaClientLogLevel("")).toBe(KafkaJS.logLevel.ERROR); + expect(kafkaClientLogLevel("not-a-level")).toBe(KafkaJS.logLevel.ERROR); +}); + +test("client log level accepts configured values regardless of case", () => { + expect(kafkaClientLogLevel("debug")).toBe(KafkaJS.logLevel.DEBUG); + expect(kafkaClientLogLevel("DEBUG")).toBe(KafkaJS.logLevel.DEBUG); + expect(kafkaClientLogLevel("Warn")).toBe(KafkaJS.logLevel.WARN); + expect(kafkaClientLogLevel("nothing")).toBe(KafkaJS.logLevel.NOTHING); +}); + +test("librdkafka debug facilities are only set when configured", () => { + expect(buildKafkaConfig(opts)).not.toHaveProperty("debug"); + expect(buildKafkaConfig(opts, {})).not.toHaveProperty("debug"); + expect(buildKafkaConfig(opts, { debug: "cgrp,fetch,broker" })).toMatchObject({ debug: "cgrp,fetch,broker" }); +}); + +test("defaults keep the client quiet", () => { + const config = buildKafkaConfig(opts); + expect(config.kafkaJS.logLevel).toBe(KafkaJS.logLevel.ERROR); + expect(config.kafkaJS.clientId).toBe("test"); + expect(config.kafkaJS.brokers).toEqual(["localhost:9092"]); +}); + +test("client logs are routed to a logCreator so librdkafka output is not discarded", () => { + const config = buildKafkaConfig(opts, { logLevel: "debug" }); + expect(config.kafkaJS.logLevel).toBe(KafkaJS.logLevel.DEBUG); + expect(typeof config.kafkaJS.logCreator).toBe("function"); + expect(() => + config.kafkaJS.logCreator()({ + level: KafkaJS.logLevel.DEBUG, + log: { fac: "CGRP", message: "rebalance: assigned 3 partition(s)" }, + }) + ).not.toThrow(); +}); diff --git a/services/rotor/src/lib/kafka-config.ts b/services/rotor/src/lib/kafka-config.ts index 3e1061ad6..c57c0d0b4 100644 --- a/services/rotor/src/lib/kafka-config.ts +++ b/services/rotor/src/lib/kafka-config.ts @@ -57,28 +57,54 @@ export function getCredentialsFromEnv(): KafkaCredentials { }; } -export function connectToKafka(opts: { defaultAppId: string } & KafkaCredentials): KafkaJS.Kafka { +const clientLogLevels: Record = { + nothing: KafkaJS.logLevel.NOTHING, + error: KafkaJS.logLevel.ERROR, + warn: KafkaJS.logLevel.WARN, + info: KafkaJS.logLevel.INFO, + debug: KafkaJS.logLevel.DEBUG, +}; + +export function kafkaClientLogLevel(configured?: string): KafkaJS.logLevel { + return clientLogLevels[(configured || "").toLowerCase()] ?? KafkaJS.logLevel.ERROR; +} + +export type KafkaDiagnostics = { + debug?: string; + logLevel?: string; +}; + +export function buildKafkaConfig( + opts: { defaultAppId: string } & KafkaCredentials, + diagnostics: KafkaDiagnostics = {} +): any { const sasl = opts.sasl ? { sasl: opts.sasl as any, } : {}; - log.atDebug().log("SASL config", JSON.stringify(opts.sasl)); - return new KafkaJS.Kafka({ + return { kafkaJS: { - logLevel: KafkaJS.logLevel.ERROR, - // logCreator: logLevel => log => { - // translateLevel(logLevel).log( - // `${log.namespace ? `${log.namespace} # ` : ""}${JSON.stringify(omit(log.log, "timestamp", "logger"))}` - // ); - // }, + logLevel: kafkaClientLogLevel(diagnostics.logLevel), + logCreator: () => (entry: any) => + translateLevel(entry.level).log( + `${entry.log?.fac ? `[${entry.log.fac}] ` : ""}${entry.log?.message ?? JSON.stringify(entry.log)}` + ), clientId: serverEnv.APPLICATION_ID || opts.defaultAppId, brokers: typeof opts.brokers === "string" ? (opts.brokers as string).split(",") : opts.brokers, ...(opts.ssl ? { ssl: true } : {}), ...sasl, }, + ...(diagnostics.debug ? { debug: diagnostics.debug } : {}), ...opts.ssl, - }); + }; +} + +export function connectToKafka(opts: { defaultAppId: string } & KafkaCredentials): KafkaJS.Kafka { + log.atDebug().log("SASL config", JSON.stringify(opts.sasl)); + return new KafkaJS.Kafka( + buildKafkaConfig(opts, { debug: serverEnv.KAFKA_DEBUG, logLevel: serverEnv.KAFKA_CLIENT_LOG_LEVEL }) + ); } export function destinationMessagesTopic(): string { diff --git a/services/rotor/src/serverEnv.ts b/services/rotor/src/serverEnv.ts index e5708c7f7..e67b70688 100644 --- a/services/rotor/src/serverEnv.ts +++ b/services/rotor/src/serverEnv.ts @@ -52,6 +52,8 @@ const ServerEnvSchema = z.object({ KAFKA_DESTINATIONS_MT_TOPIC_NAME: z.string().optional().default("destination-messages-mt"), KAFKA_CONSUMER_GROUP_ID: z.string().optional(), KAFKA_TOPIC_COMPRESSION: z.string().optional().default("gzip"), + KAFKA_DEBUG: z.string().optional(), + KAFKA_CLIENT_LOG_LEVEL: z.string().optional(), CONSUMER_PROTOCOL: z.string().optional(), // Bulker Configuration