Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 43 additions & 0 deletions services/rotor/__tests__/kafka-config.test.ts
Original file line number Diff line number Diff line change
@@ -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();
});
46 changes: 36 additions & 10 deletions services/rotor/src/lib/kafka-config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,28 +57,54 @@ export function getCredentialsFromEnv(): KafkaCredentials {
};
}

export function connectToKafka(opts: { defaultAppId: string } & KafkaCredentials): KafkaJS.Kafka {
const clientLogLevels: Record<string, KafkaJS.logLevel> = {
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 {
Expand Down
2 changes: 2 additions & 0 deletions services/rotor/src/serverEnv.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down