From c4d8b4a67b33a61e7a1b602b99b59abf08aad22b Mon Sep 17 00:00:00 2001 From: root Date: Sat, 18 Jul 2026 08:51:06 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20=E6=94=AF=E6=8C=81=20L1=20Kafka=20?= =?UTF-8?q?=E9=9B=86=E7=BE=A4=E5=9F=9F=E5=90=8D=E8=A7=A3=E6=9E=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/cloud/kafka-event-bridge.test.ts | 37 ++++++++- internal/cloud/kafka-event-bridge.ts | 99 ++++++++++++++++++++++- tools/hwlab-cli/kafka-regenerate.test.ts | 6 +- tools/src/hwlab-cli/kafka-regenerate.ts | 18 ++++- 4 files changed, 149 insertions(+), 11 deletions(-) diff --git a/internal/cloud/kafka-event-bridge.test.ts b/internal/cloud/kafka-event-bridge.test.ts index 907206a0..dc0b0753 100644 --- a/internal/cloud/kafka-event-bridge.test.ts +++ b/internal/cloud/kafka-event-bridge.test.ts @@ -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", diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 56c7e60a..3b360673 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -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) { diff --git a/tools/hwlab-cli/kafka-regenerate.test.ts b/tools/hwlab-cli/kafka-regenerate.test.ts index 397e2ea2..efb720ee 100644 --- a/tools/hwlab-cli/kafka-regenerate.test.ts +++ b/tools/hwlab-cli/kafka-regenerate.test.ts @@ -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"); diff --git a/tools/src/hwlab-cli/kafka-regenerate.ts b/tools/src/hwlab-cli/kafka-regenerate.ts index 3ddd6378..5ed43d16 100644 --- a/tools/src/hwlab-cli/kafka-regenerate.ts +++ b/tools/src/hwlab-cli/kafka-regenerate.ts @@ -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) { 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(task: () => Promise) { 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 } : {}) }); }