From 8767e9452b136de055716f33569ca74d4b3d117b Mon Sep 17 00:00:00 2001 From: Tyagiquamar Date: Fri, 4 Sep 2026 17:15:22 +0530 Subject: [PATCH 1/4] fix(links): fail health check when Kafka delivery is unavailable --- apps/links/src/index.ts | 4 + apps/links/src/lib/producer.test.ts | 120 ++++++++++++++++++++++++++++ apps/links/src/lib/producer.ts | 12 ++- 3 files changed, 135 insertions(+), 1 deletion(-) diff --git a/apps/links/src/index.ts b/apps/links/src/index.ts index f5afeec5f..13f40f226 100644 --- a/apps/links/src/index.ts +++ b/apps/links/src/index.ts @@ -12,6 +12,7 @@ import { evlog } from "evlog/elysia"; import { drain, enrich, flushDrain } from "./lib/logging"; import { calculateLinkReadiness } from "./lib/health"; import { + didKafkaConnectFail, disconnectProducer, getProducerHealthState, refreshProducerConnection, @@ -200,6 +201,9 @@ const app = new Elysia() redis: cache.status, redpanda: redpanda.status, }); + if (didKafkaConnectFail()) { + return Response.json({ status: "unavailable", services }, { status: 503 }); + } return Response.json( { status: readiness.status, diff --git a/apps/links/src/lib/producer.test.ts b/apps/links/src/lib/producer.test.ts index 511915442..049e8ff34 100644 --- a/apps/links/src/lib/producer.test.ts +++ b/apps/links/src/lib/producer.test.ts @@ -308,3 +308,123 @@ describe("producer health state", () => { } }); }); + +describe("health probe failures (issue #719)", () => { + test("health refresh rejects loudly when Kafka is unreachable", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + nextProducer = makeProducer({ + connect: () => Promise.reject(new Error("tls handshake failed")), + }); + const { getProducerHealthState, refreshProducerConnection } = + await loadProducer(); + + await expect(refreshProducerConnection()).rejects.toThrow( + "Kafka health probe failed" + ); + expect(setAttributes).toHaveBeenCalledWith({ + kafka_health_connect_failed: true, + }); + expect(getProducerHealthState()).toBe("cooldown"); + }); + + test("health refresh resolves once Kafka is reachable", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + nextProducer = makeProducer(); + const { + disconnectProducer, + getProducerHealthState, + refreshProducerConnection, + } = await loadProducer(); + + await refreshProducerConnection(); + + expect(getProducerHealthState()).toBe("connected"); + await disconnectProducer(); + }); + + test("production sends still fall back to ClickHouse when Kafka is down", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + nextProducer = makeProducer({ + connect: () => Promise.reject(new Error("broker unavailable")), + }); + const { sendLinkVisit } = await loadProducer(); + + const result = await sendLinkVisit(event, event.link_id); + + expect(result).toBe(true); + expect(clickHouseInsert).toHaveBeenCalledTimes(1); + }); + + test("flags Kafka outages for the health endpoint until recovery", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + const realNow = Date.now; + let now = 1000; + Date.now = () => now; + try { + nextProducer = makeProducer({ + connect: () => Promise.reject(new Error("broker unavailable")), + }); + const { + didKafkaConnectFail, + disconnectProducer, + refreshProducerConnection, + warmProducerConnection, + } = await loadProducer(); + + expect(didKafkaConnectFail()).toBe(false); + await warmProducerConnection(); + expect(didKafkaConnectFail()).toBe(true); + + now += 60_001; + nextProducer = makeProducer(); + await refreshProducerConnection(); + + expect(didKafkaConnectFail()).toBe(false); + await disconnectProducer(); + } finally { + Date.now = realNow; + } + }); + + test("suppressed health probes still report a known outage", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + nextProducer = makeProducer({ + connect: () => Promise.reject(new Error("tls handshake failed")), + }); + const { didKafkaConnectFail, refreshProducerConnection } = + await loadProducer(); + + await expect(refreshProducerConnection()).rejects.toThrow( + "Kafka health probe failed" + ); + expect(didKafkaConnectFail()).toBe(true); + + await expect(refreshProducerConnection()).rejects.toThrow( + "Kafka health probe failed" + ); + expect(didKafkaConnectFail()).toBe(true); + }); + + test("a failing health-owned attempt does not break concurrent sends", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + let releaseConnect!: (error: Error) => void; + nextProducer = makeProducer({ + connect: () => + new Promise((_resolve, reject) => { + releaseConnect = reject; + }), + }); + const { refreshProducerConnection, sendLinkVisit } = + await loadProducer(); + + const health = refreshProducerConnection(); + await Bun.sleep(0); + const send = sendLinkVisit(event, event.link_id); + await Bun.sleep(0); + releaseConnect(new Error("tls handshake failed")); + + await expect(health).rejects.toThrow("Kafka health probe failed"); + await expect(send).resolves.toBe(true); + expect(clickHouseInsert).toHaveBeenCalledTimes(1); + }); +}); diff --git a/apps/links/src/lib/producer.ts b/apps/links/src/lib/producer.ts index ae5b1d99b..7900a686c 100644 --- a/apps/links/src/lib/producer.ts +++ b/apps/links/src/lib/producer.ts @@ -16,6 +16,7 @@ const ASYNC_INSERT_BUSY_TIMEOUT_MS = 50; let producer: Producer | null = null; let connectPromise: Promise | null = null; let nextReconnectAt = 0; +let kafkaConnectFailed = false; let lastConnectErrorLogAt = 0; let lastFallbackErrorLogAt = 0; let shuttingDown = false; @@ -131,6 +132,7 @@ function connect(reportFailure = true): Promise { } producer = candidate; nextReconnectAt = 0; + kafkaConnectFailed = false; setAttributes({ kafka_connected: true }); return true; } catch (error) { @@ -149,6 +151,7 @@ function connect(reportFailure = true): Promise { producer = null; nextReconnectAt = Date.now() + reconnectCooldownMs; setAttributes({ kafka_connected: false }); + kafkaConnectFailed = true; return false; } finally { connectPromise = null; @@ -181,12 +184,19 @@ export function getProducerHealthState(): ProducerHealthState { return "idle"; } +export function didKafkaConnectFail(): boolean { + return kafkaConnectFailed; +} + export async function warmProducerConnection(): Promise { await connect(); } export async function refreshProducerConnection(): Promise { - await connect(false); + const connected = await connect(false); + if (!connected) { + throw new Error("Kafka health probe failed"); + } } async function persistLinkVisitDirectly( From 27a313b9e0117a6e3a6ac836a2ce6e9f49c6421f Mon Sep 17 00:00:00 2001 From: Tyagiquamar Date: Sat, 5 Sep 2026 10:19:35 +0530 Subject: [PATCH 2/4] fix(links): handle health probe promise sharing, TLS config and shutdown safety --- .github/workflows/health-check.yml | 3 +- apps/links/src/lib/producer.test.ts | 64 +++++++++++++++++++++++++++++ apps/links/src/lib/producer.ts | 31 ++++++++++++-- 3 files changed, 93 insertions(+), 5 deletions(-) diff --git a/.github/workflows/health-check.yml b/.github/workflows/health-check.yml index 728b3c13e..1ddfdc2ec 100644 --- a/.github/workflows/health-check.yml +++ b/.github/workflows/health-check.yml @@ -739,6 +739,7 @@ jobs: -e REDIS_URL=redis://localhost:6379 \ -e BULLMQ_REDIS_URL=redis://localhost:6379/4 \ -e REDPANDA_BROKER=localhost:9092 \ + -e REDPANDA_SSL=false \ -e DASHBOARD_URL=http://localhost:3000 \ -e CLICKHOUSE_URL=http://default:@localhost:8123/databuddy_analytics \ links:test @@ -757,7 +758,7 @@ jobs: for i in {1..30}; do STATUS_BODY=$(curl -sS http://localhost:2500/health/status) echo "Links /health/status: $STATUS_BODY" - if echo "$STATUS_BODY" | jq -e '.status == "ok" or .status == "degraded"' > /dev/null; then + if echo "$STATUS_BODY" | jq -e '.status == "ok"' > /dev/null; then echo "Links dependency health is valid" break fi diff --git a/apps/links/src/lib/producer.test.ts b/apps/links/src/lib/producer.test.ts index 049e8ff34..926833462 100644 --- a/apps/links/src/lib/producer.test.ts +++ b/apps/links/src/lib/producer.test.ts @@ -427,4 +427,68 @@ describe("health probe failures (issue #719)", () => { await expect(send).resolves.toBe(true); expect(clickHouseInsert).toHaveBeenCalledTimes(1); }); + + test("disables SSL when REDPANDA_SSL is false", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + process.env.REDPANDA_SSL = "false"; + nextProducer = makeProducer(); + const { disconnectProducer, sendLinkVisit } = await loadProducer(); + + await sendLinkVisit(event, event.link_id); + + expect(kafkaConfigs).toEqual([ + expect.objectContaining({ + brokers: ["redpanda.test:9092"], + ssl: false, + }), + ]); + await disconnectProducer(); + }); + + test("reports connection dependency error when sendLinkVisit joins health-owned attempt", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + let releaseConnect!: (error: Error) => void; + const connectErr = new Error("connection failed"); + nextProducer = makeProducer({ + connect: () => + new Promise((_resolve, reject) => { + releaseConnect = reject; + }), + }); + const { refreshProducerConnection, sendLinkVisit } = + await loadProducer(); + + const health = refreshProducerConnection(); + await Bun.sleep(0); + const send = sendLinkVisit(event, event.link_id); + await Bun.sleep(0); + releaseConnect(connectErr); + + await expect(health).rejects.toThrow("Kafka health probe failed"); + await expect(send).resolves.toBe(true); + expect(captureError).toHaveBeenCalledWith(connectErr, { + operation: "kafka_connect", + }); + }); + + test("does not fail disconnectProducer if a pending connection fails during shutdown", async () => { + process.env.REDPANDA_BROKER = "redpanda.test:9092"; + let releaseConnect!: (error: Error) => void; + nextProducer = makeProducer({ + connect: () => + new Promise((_resolve, reject) => { + releaseConnect = reject; + }), + }); + const { disconnectProducer, warmProducerConnection } = + await loadProducer(); + + const warmup = warmProducerConnection(); + await Bun.sleep(0); + const shutdown = disconnectProducer(); + releaseConnect(new Error("connection rejected during shutdown")); + + await warmup; + await expect(shutdown).resolves.toBeUndefined(); + }); }); diff --git a/apps/links/src/lib/producer.ts b/apps/links/src/lib/producer.ts index 7900a686c..1f3130614 100644 --- a/apps/links/src/lib/producer.ts +++ b/apps/links/src/lib/producer.ts @@ -70,6 +70,8 @@ function captureDependencyError( captureError(error, context); } +let shouldReportFailure = false; + function connect(reportFailure = true): Promise { if (shuttingDown) { return Promise.resolve(false); @@ -82,6 +84,10 @@ function connect(reportFailure = true): Promise { return Promise.resolve(false); } + if (reportFailure) { + shouldReportFailure = true; + } + const now = Date.now(); if (now < nextReconnectAt) { setAttributes({ kafka_reconnect_suppressed: true }); @@ -92,6 +98,12 @@ function connect(reportFailure = true): Promise { return connectPromise; } + const useSsl = + process.env.REDPANDA_SSL === "false" || + process.env.REDPANDA_SSL_ENABLED === "false" + ? false + : true; + connectPromise = (async () => { let candidate: Producer | null = null; @@ -111,7 +123,7 @@ function connect(reportFailure = true): Promise { ...(username && password ? { sasl: { mechanism: "scram-sha-256", username, password } } : {}), - ssl: true, + ssl: useSsl, }); candidate = kafka.producer({ @@ -136,7 +148,7 @@ function connect(reportFailure = true): Promise { setAttributes({ kafka_connected: true }); return true; } catch (error) { - if (reportFailure) { + if (shouldReportFailure) { captureDependencyError( error, { operation: "kafka_connect" }, @@ -155,6 +167,7 @@ function connect(reportFailure = true): Promise { return false; } finally { connectPromise = null; + shouldReportFailure = false; } })(); @@ -259,7 +272,13 @@ export async function sendLinkVisit( kafka_message_key: eventKey ?? "unknown", }); - const kafkaReady = await connect(); + let kafkaReady = false; + try { + kafkaReady = await connect(); + } catch (error) { + kafkaReady = false; + } + const activeProducer = producer; if (!(kafkaReady && activeProducer)) { setAttributes({ @@ -304,7 +323,11 @@ export async function sendLinkVisit( export async function disconnectProducer(): Promise { shuttingDown = true; if (connectPromise) { - await connectPromise; + try { + await connectPromise; + } catch { + // Expected health/connect probe rejection during shutdown should not fail disconnectProducer + } } const activeProducer = producer; producer = null; From 12898bd3f5564814a24ca97440ae4324393fddee Mon Sep 17 00:00:00 2001 From: Tyagiquamar Date: Sat, 5 Sep 2026 10:49:03 +0530 Subject: [PATCH 3/4] fix(links): un-modify workflow file to pass Tripwire firewall check --- .github/workflows/health-check.yml | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/.github/workflows/health-check.yml b/.github/workflows/health-check.yml index 1ddfdc2ec..728b3c13e 100644 --- a/.github/workflows/health-check.yml +++ b/.github/workflows/health-check.yml @@ -739,7 +739,6 @@ jobs: -e REDIS_URL=redis://localhost:6379 \ -e BULLMQ_REDIS_URL=redis://localhost:6379/4 \ -e REDPANDA_BROKER=localhost:9092 \ - -e REDPANDA_SSL=false \ -e DASHBOARD_URL=http://localhost:3000 \ -e CLICKHOUSE_URL=http://default:@localhost:8123/databuddy_analytics \ links:test @@ -758,7 +757,7 @@ jobs: for i in {1..30}; do STATUS_BODY=$(curl -sS http://localhost:2500/health/status) echo "Links /health/status: $STATUS_BODY" - if echo "$STATUS_BODY" | jq -e '.status == "ok"' > /dev/null; then + if echo "$STATUS_BODY" | jq -e '.status == "ok" or .status == "degraded"' > /dev/null; then echo "Links dependency health is valid" break fi From 62f865de888a7c05b39ef9d210029a9183edfafa Mon Sep 17 00:00:00 2001 From: Tyagiquamar Date: Sat, 5 Sep 2026 22:57:05 +0530 Subject: [PATCH 4/4] fix(links): set kafka failure flag only after reconnect cooldown check --- apps/links/src/lib/producer.ts | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/apps/links/src/lib/producer.ts b/apps/links/src/lib/producer.ts index 1f3130614..183399ea7 100644 --- a/apps/links/src/lib/producer.ts +++ b/apps/links/src/lib/producer.ts @@ -84,16 +84,19 @@ function connect(reportFailure = true): Promise { return Promise.resolve(false); } - if (reportFailure) { - shouldReportFailure = true; - } - const now = Date.now(); if (now < nextReconnectAt) { setAttributes({ kafka_reconnect_suppressed: true }); return Promise.resolve(false); } + // Set only after the cooldown early return: that return path skips the + // finally that resets the flag, so setting it earlier would leak a stale + // reportFailure=true into the next real connect attempt. + if (reportFailure) { + shouldReportFailure = true; + } + if (connectPromise) { return connectPromise; }