fix: 支持 L1 Kafka 集群域名解析

This commit is contained in:
root
2026-07-18 08:51:06 +02:00
parent f1597e2b6b
commit c4d8b4a67b
4 changed files with 149 additions and 11 deletions
+36 -1
View File
@@ -6,7 +6,7 @@ import { test } from "bun:test";
import { createCloudApiBunServer } from "./bun-server.ts";
import { buildCloudApiReadiness } from "./health-contract.ts";
import { decodeCanonicalAgentRunKafkaMessage, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
import { decodeCanonicalAgentRunKafkaMessage, kafkaDnsLookup, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
const PROJECTOR_ENV = Object.freeze({
HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false",
@@ -58,6 +58,41 @@ const LIVE_REFRESH_ENV = Object.freeze({
HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT: "2000"
});
test("uses explicitly configured DNS servers for Kafka broker metadata hostnames", async () => {
class FakeResolver {
static instance;
constructor() { FakeResolver.instance = this; }
setServers(servers) { this.servers = servers; }
resolve4(hostname, callback) {
this.hostnames ??= [];
this.hostnames.push(hostname);
if (hostname.endsWith(".cluster.local")) callback(null, ["10.42.0.64"]);
else callback(Object.assign(new Error("not found"), { code: "ENOTFOUND" }));
}
}
const lookup = kafkaDnsLookup(["10.43.0.10"], ["cluster.local"], FakeResolver);
const resolved = await new Promise((resolve, reject) => {
lookup("platform-infra-kafka-dual-role-0.platform-infra-kafka-kafka-brokers.platform-infra.svc", {}, (error, address, family) => {
if (error) reject(error);
else resolve({ address, family });
});
});
assert.deepEqual(FakeResolver.instance.servers, ["10.43.0.10"]);
assert.deepEqual(FakeResolver.instance.hostnames, [
"platform-infra-kafka-dual-role-0.platform-infra-kafka-kafka-brokers.platform-infra.svc",
"platform-infra-kafka-dual-role-0.platform-infra-kafka-kafka-brokers.platform-infra.svc.cluster.local"
]);
assert.deepEqual(resolved, { address: "10.42.0.64", family: 4 });
assert.equal(kafkaDnsLookup([]), undefined);
const config = kafkaEventBridgeConfig({
...LIVE_ENV,
HWLAB_KAFKA_DNS_SERVERS: "10.43.0.10,10.43.0.11",
HWLAB_KAFKA_DNS_SEARCH_DOMAINS: "cluster.local"
});
assert.deepEqual(config.dnsServers, ["10.43.0.10", "10.43.0.11"]);
assert.deepEqual(config.dnsSearchDomains, ["cluster.local"]);
});
test("projects only explicit AgentRun lifecycle timestamps", () => {
const projected = projectAgentRunKafkaEventToHwlabEvent({
schema: "agentrun.event.v1",
+95 -4
View File
@@ -2,6 +2,10 @@
// Responsibility: subscribe AgentRun Kafka events, project them into HWLAB Kafka events, and provide bounded Kafka stream queries for diagnostics.
import { createHash, randomUUID } from "node:crypto";
import { Resolver } from "node:dns";
import { isIP } from "node:net";
import net from "node:net";
import tls from "node:tls";
import { Kafka, logLevel } from "kafkajs";
import { emitCodeAgentOtelSpan } from "./otel-trace.ts";
@@ -68,6 +72,8 @@ export function kafkaEventBridgeConfig(env = process.env) {
throw error;
}
const brokers = csv(env.HWLAB_KAFKA_BOOTSTRAP_SERVERS);
const dnsServers = csv(env.HWLAB_KAFKA_DNS_SERVERS);
const dnsSearchDomains = csv(env.HWLAB_KAFKA_DNS_SEARCH_DOMAINS);
const agentRunTopic = stringValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC);
const hwlabTopic = stringValue(env.HWLAB_KAFKA_EVENT_TOPIC);
const clientId = stringValue(env.HWLAB_KAFKA_CLIENT_ID);
@@ -104,6 +110,8 @@ export function kafkaEventBridgeConfig(env = process.env) {
return {
capabilities,
brokers,
dnsServers,
dnsSearchDomains,
agentRunTopic,
hwlabTopic,
clientId,
@@ -961,7 +969,12 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
: strictPositiveInteger(scanLimit, "scanLimit");
const budgetMs = Math.max(250, integerValue(timeoutMs) || DEFAULT_QUERY_TIMEOUT_MS);
const deadlineMs = Date.now() + budgetMs;
const kafka = kafkaFactory({ brokers, clientId: `${clientId}-query` });
const kafka = kafkaFactory({
brokers,
clientId: `${clientId}-query`,
dnsServers: csv(env.HWLAB_KAFKA_DNS_SERVERS),
dnsSearchDomains: csv(env.HWLAB_KAFKA_DNS_SEARCH_DOMAINS)
});
const endOffsetSnapshot = await kafkaTopicEndOffsetSnapshot(kafka, resolvedTopic, deadlineMs, signal);
const endOffsetByPartition = new Map(endOffsetSnapshot.offsets.map((entry) => [entry.partition, entry.endOffset]));
const startOffsetByPartition = new Map(endOffsetSnapshot.offsets.map((entry) => [entry.partition, entry.startOffset]));
@@ -1210,7 +1223,12 @@ export async function openKafkaEventStream({ env = process.env, stream = "hwlab"
const resolvedTopic = stringValue(topic) || kafkaTopicForStream(stream, env);
const filters = compactObject({ traceId, sessionId, runId, commandId, replayId });
const resolvedOffsetRange = normalizeKafkaOffsetRange(offsetRange);
const kafka = kafkaFactory({ brokers, clientId: `${clientId}-debug-sse` });
const kafka = kafkaFactory({
brokers,
clientId: `${clientId}-debug-sse`,
dnsServers: csv(env.HWLAB_KAFKA_DNS_SERVERS),
dnsSearchDomains: csv(env.HWLAB_KAFKA_DNS_SEARCH_DOMAINS)
});
const resolvedGroupIdPrefix = stringValue(groupIdPrefix) || `${clientId}-debug-sse`;
const groupId = `${resolvedGroupIdPrefix}-${Date.now()}-${randomUUID().slice(0, 8)}`;
const consumer = kafka.consumer({ groupId, allowAutoTopicCreation: false });
@@ -1564,8 +1582,81 @@ function kafkaTopicForStream(stream, env) {
return normalized;
}
function defaultKafkaFactory(config) {
return new Kafka({ clientId: config.clientId, brokers: config.brokers, logLevel: logLevel.NOTHING });
export function kafkaDnsLookup(dnsServers, dnsSearchDomains = [], ResolverClass = Resolver) {
if (!Array.isArray(dnsServers) || dnsServers.length === 0) return undefined;
const resolver = new ResolverClass();
resolver.setServers(dnsServers);
return (hostname, options, callback) => {
const family = typeof options === "number" ? options : options?.family;
const all = typeof options === "object" && options?.all === true;
const literalFamily = isIP(hostname);
if (literalFamily > 0) {
callback(null, all ? [{ address: hostname, family: literalFamily }] : hostname, literalFamily);
return;
}
if (family === 6) {
resolveKafkaHostname(resolver, hostname, dnsSearchDomains, 6, (error, addresses) => finishKafkaDnsLookup(error, addresses, 6, all, callback));
return;
}
resolveKafkaHostname(resolver, hostname, dnsSearchDomains, 4, (error, addresses) => finishKafkaDnsLookup(error, addresses, 4, all, callback));
};
}
function resolveKafkaHostname(resolver, hostname, searchDomains, family, callback) {
const candidates = [
hostname,
...searchDomains
.map((domain) => domain.replace(/^\.+|\.+$/gu, ""))
.filter(Boolean)
.filter((domain) => !hostname.endsWith(`.${domain}`))
.map((domain) => `${hostname}.${domain}`)
];
const resolveNext = (index, previousError = null) => {
if (index >= candidates.length) {
callback(previousError);
return;
}
resolver[family === 6 ? "resolve6" : "resolve4"](candidates[index], (error, addresses) => {
if (!error && Array.isArray(addresses) && addresses.length > 0) {
callback(null, addresses);
return;
}
resolveNext(index + 1, error ?? previousError);
});
};
resolveNext(0);
}
export function kafkaSocketFactory(dnsServers, dnsSearchDomains = []) {
const lookup = kafkaDnsLookup(dnsServers, dnsSearchDomains);
if (!lookup) return undefined;
return ({ host, port, ssl, onConnect }) => ssl
? tls.connect({ host, port, ...ssl, lookup }, onConnect)
: net.connect({ host, port, lookup }, onConnect);
}
function finishKafkaDnsLookup(error, addresses, family, all, callback) {
if (error) {
callback(error);
return;
}
if (!Array.isArray(addresses) || addresses.length === 0) {
const emptyError = new Error(`Kafka DNS lookup returned no addresses for family IPv${family}.`);
emptyError.code = "hwlab_kafka_dns_empty";
callback(emptyError);
return;
}
callback(null, all ? addresses.map((address) => ({ address, family })) : addresses[0], family);
}
export function defaultKafkaFactory(config) {
const socketFactory = kafkaSocketFactory(config.dnsServers, config.dnsSearchDomains);
return new Kafka({
clientId: config.clientId,
brokers: config.brokers,
logLevel: logLevel.NOTHING,
...(socketFactory ? { socketFactory } : {})
});
}
async function boundedDisconnect(consumer, timeoutMs = 2000) {
+4 -2
View File
@@ -827,8 +827,9 @@ test("Kafka Trace render returns a typed partial result when the captured end of
filterRejectedCount: 0,
completionReason: "timeout",
reachedEndOffsets: false,
endOffsetsAvailable: true,
endOffsets: [{ partition: 0, startOffset: "0", endOffset: "2" }],
endOffsetsAvailable: false,
endOffsetsError: { code: "ENOTFOUND", message: "broker metadata hostname unresolved", valuesRedacted: true },
endOffsets: [],
completion: { reason: "timeout", complete: false, barrierReached: false, retentionStartVerified: true },
lastScannedOffsets: [{ partition: 0, offset: "0" }]
};
@@ -836,6 +837,7 @@ test("Kafka Trace render returns a typed partial result when the captured end of
});
assert.equal(result.exitCode, 2);
assert.equal(result.payload.input.endOffsetsError.code, "ENOTFOUND");
assert.equal(result.payload.ok, false);
assert.equal(result.payload.status, "partial");
assert.equal(result.payload.error.code, "source_scan_incomplete");
+14 -4
View File
@@ -12,6 +12,7 @@ import {
DEFAULT_AGENTRUN_EVENT_TOPIC,
DEFAULT_HWLAB_DEBUG_EVENT_TOPIC,
DEFAULT_HWLAB_EVENT_TOPIC,
kafkaSocketFactory,
projectAgentRunKafkaEventToHwlabEvent,
projectAgentRunKafkaMessageToHwlabDebugEvent,
queryKafkaEventStream
@@ -427,6 +428,7 @@ export async function renderHwlabKafkaTrace(parsed: ParsedArgs, dependencies: Re
completionReason: readResult.completionReason ?? null,
reachedEndOffsets: readResult.reachedEndOffsets ?? false,
endOffsetsAvailable: readResult.endOffsetsAvailable ?? false,
endOffsetsError: readResult.endOffsetsError ?? null,
endOffsets: readResult.endOffsets ?? [],
lastScannedOffsets: readResult.lastScannedOffsets ?? [],
sourceFile: readResult.sourceFile ?? null,
@@ -598,6 +600,7 @@ export async function renderHwlabKafkaTrace(parsed: ParsedArgs, dependencies: Re
completionReason: readResult.completionReason ?? (sourceMode === "jsonl" ? "file-end" : null),
reachedEndOffsets: readResult.reachedEndOffsets ?? (sourceMode === "jsonl"),
endOffsetsAvailable: readResult.endOffsetsAvailable ?? false,
endOffsetsError: readResult.endOffsetsError ?? null,
endOffsets: readResult.endOffsets ?? [],
lastScannedOffsets: readResult.lastScannedOffsets ?? [],
sourceFile: readResult.sourceFile ?? null,
@@ -1661,7 +1664,7 @@ function kafkaHelp() {
async function defaultKafkaReader(input: Record<string, any>) {
return withBunKafkaTimerClamp(() => queryKafkaEventStream({
...input,
kafkaFactory(config: { brokers: string[]; clientId: string }) { return boundedKafka(config); }
kafkaFactory(config: { brokers: string[]; clientId: string; dnsServers?: string[]; dnsSearchDomains?: string[] }) { return boundedKafka(config); }
}));
}
@@ -1685,17 +1688,24 @@ async function withBunKafkaTimerClamp<T>(task: () => Promise<T>) {
async function defaultProducerFactory({ env, clientId }: { env: EnvLike; clientId: string }) {
const brokers = csv(env.HWLAB_KAFKA_BOOTSTRAP_SERVERS);
if (brokers.length === 0) throw cliError("kafka_brokers_missing", "HWLAB_KAFKA_BOOTSTRAP_SERVERS is required for Kafka publish", { field: "HWLAB_KAFKA_BOOTSTRAP_SERVERS" });
return boundedKafka({ clientId, brokers }).producer({ allowAutoTopicCreation: false });
return boundedKafka({
clientId,
brokers,
dnsServers: csv(env.HWLAB_KAFKA_DNS_SERVERS),
dnsSearchDomains: csv(env.HWLAB_KAFKA_DNS_SEARCH_DOMAINS)
}).producer({ allowAutoTopicCreation: false });
}
function boundedKafka({ clientId, brokers }: { clientId: string; brokers: string[] }) {
function boundedKafka({ clientId, brokers, dnsServers = [], dnsSearchDomains = [] }: { clientId: string; brokers: string[]; dnsServers?: string[]; dnsSearchDomains?: string[] }) {
const socketFactory = kafkaSocketFactory(dnsServers, dnsSearchDomains);
return new Kafka({
clientId,
brokers,
logLevel: logLevel.NOTHING,
connectionTimeout: 5000,
requestTimeout: 10000,
retry: { retries: 1, maxRetryTime: 10000 }
retry: { retries: 1, maxRetryTime: 10000 },
...(socketFactory ? { socketFactory } : {})
});
}