Skip to content
Merged
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
1 change: 1 addition & 0 deletions .agents/skills/databuddy-internal/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -187,6 +187,7 @@ Read [codebase-map.md](./references/codebase-map.md) when you need deeper routin
- `pg.Pool` already grows lazily from zero to its configured `max`; do not replace it with one `Client` to address acquisition timeouts, because that serializes queries. Keep a bounded pool, tune its acquisition timeout, and monitor `waitingCount`.
- ClickHouse helpers and schema: `packages/db/src/clickhouse/*`
- After schema changes, use the repo db scripts rather than ad hoc commands
- Do not add ClickHouse migration files for delivery hardening; keep relay identity in the worker/queue unless **iza** explicitly requests persistent warehouse identity.

### Auth and permissions

Expand Down
3 changes: 2 additions & 1 deletion apps/uptime/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,8 @@
"effect": "^4.0.0-beta.59",
"elysia": "catalog:",
"evlog": "catalog:",
"kafkajs": "^2.2.4"
"kafkajs": "^2.2.4",
"zod": "catalog:"
},
"packageManager": "bun@1.3.14"
}
2 changes: 2 additions & 0 deletions apps/uptime/src/actions.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { randomUUID } from "node:crypto";
import { connect } from "node:tls";
import { db } from "@databuddy/db";
import {
Expand Down Expand Up @@ -431,6 +432,7 @@ const runUptimeCheck = (
site_id: siteId,
url: normalizedUrl,
timestamp,
event_id: randomUUID(),
status: pingResult.ok ? MonitorStatus.UP : MonitorStatus.DOWN,
http_code: pingResult.statusCode,
ttfb_ms: pingResult.ttfb,
Expand Down
118 changes: 88 additions & 30 deletions apps/uptime/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ import {
import { disconnectProducer } from "./lib/producer";
import { captureError } from "./lib/tracing";
import { syncSchedulers } from "./sync-schedulers";
import { startUptimeWorker } from "./worker";
import { startUptimeDeliveryWorker, startUptimeWorker } from "./worker";

initLogger({
env: createDatabuddyEvlogEnv("uptime"),
Expand All @@ -27,6 +27,12 @@ initLogger({
sampling: {},
});

let shuttingDown = false;
let shutdownExitCode = 0;
let uptimeWorker: ReturnType<typeof startUptimeWorker> | null = null;
let uptimeDeliveryWorker: ReturnType<typeof startUptimeDeliveryWorker> | null =
null;

process.on("unhandledRejection", (reason, _promise) => {
captureError(reason, { process: "unhandledRejection" });
log.error({
Expand All @@ -43,54 +49,106 @@ process.on("uncaughtException", (error) => {
error_stack: error instanceof Error ? error.stack : undefined,
error_source: "process",
});
shutdown("uncaughtException", 1).catch((shutdownError) => {
captureError(shutdownError, {
process: "uncaughtException",
error_step: "fatal_shutdown",
});
process.exit(1);
});
});

const DRAIN_TIMEOUT_MS = 10_000;

const drainAll = (worker: ReturnType<typeof startUptimeWorker> | null) =>
Effect.all(
[
Effect.tryPromise({
try: () => worker?.close() ?? Promise.resolve(),
catch: (c) => c,
}),
Effect.tryPromise({ try: () => closeUptimeQueue(), catch: (c) => c }),
Effect.tryPromise({
try: () => flushBatchedUptimeDrain(),
catch: (c) => c,
}),
Effect.tryPromise({
try: () => shutdownPostgres(),
catch: (c) => c,
}),
Effect.tryPromise({ try: () => disconnectProducer(), catch: (c) => c }),
],
{ concurrency: "unbounded" }
).pipe(
const drainStep = (step: string, action: () => Promise<void>) =>
Effect.tryPromise({
try: action,
catch: (cause) => cause,
}).pipe(
Effect.catch((cause) =>
Effect.sync(() =>
log.error({
lifecycle: "shutdown",
error_step: step,
error_message: cause instanceof Error ? cause.message : String(cause),
})
)
)
);

const drainAll = (
worker: ReturnType<typeof startUptimeWorker> | null,
deliveryWorker: ReturnType<typeof startUptimeDeliveryWorker> | null
) =>
Effect.gen(function* () {
// Stop source admission before closing the relay, preserving queued events
// for the next worker process if the shutdown window expires.
yield* drainStep(
"uptime_worker_close",
() => worker?.close() ?? Promise.resolve()
);
yield* drainStep(
"uptime_delivery_worker_close",
() => deliveryWorker?.close() ?? Promise.resolve()
);
yield* Effect.all(
[
drainStep("uptime_queue_close", () => closeUptimeQueue()),
drainStep("uptime_log_flush", () => flushBatchedUptimeDrain()),
drainStep("uptime_postgres_close", () => shutdownPostgres()),
drainStep("uptime_producer_disconnect", () => disconnectProducer()),
],
{ concurrency: "unbounded" }
);
}).pipe(
Effect.timeout(`${DRAIN_TIMEOUT_MS} millis`),
Effect.catch(() =>
Effect.catch((cause) =>
Effect.sync(() =>
log.error({
lifecycle: "shutdown",
error_step: "drain_timeout",
error_step:
cause &&
typeof cause === "object" &&
"_tag" in cause &&
cause._tag === "TimeoutError"
? "drain_timeout"
: "drain_failed",
drain_timeout_ms: DRAIN_TIMEOUT_MS,
error_message: cause instanceof Error ? cause.message : String(cause),
})
)
)
);

async function shutdown(signal: string) {
async function shutdown(signal: string, exitCode = 0) {
shutdownExitCode = Math.max(shutdownExitCode, exitCode);
if (shuttingDown) {
return;
}
shuttingDown = true;
log.info("lifecycle", `${signal} received, shutting down gracefully`);
await Effect.runPromise(drainAll(uptimeWorker));
process.exit(0);
try {
await Effect.runPromise(drainAll(uptimeWorker, uptimeDeliveryWorker));
} finally {
process.exit(shutdownExitCode);
}
}

let uptimeWorker: ReturnType<typeof startUptimeWorker> | null = null;

(async () => {
if (UPTIME_ENV.isProduction) {
await syncSchedulers();
uptimeWorker = startUptimeWorker();
try {
await syncSchedulers();
uptimeDeliveryWorker = startUptimeDeliveryWorker();
uptimeWorker = startUptimeWorker();
} catch (error) {
captureError(error, { error_step: "uptime_startup" });
log.error({
lifecycle: "startup",
error_step: "uptime_startup",
error_message: error instanceof Error ? error.message : String(error),
});
await shutdown("startup", 1);
}
} else {
log.info(
"lifecycle",
Expand Down
155 changes: 155 additions & 0 deletions apps/uptime/src/lib/producer.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
import { afterAll, beforeEach, describe, expect, mock, test } from "bun:test";

const kafkaConfigs: unknown[] = [];
const producers: Array<ReturnType<typeof createProducer>> = [];
const captureError = mock(() => {});

class KafkaMock {
constructor(config: unknown) {
kafkaConfigs.push(config);
}

producer() {
const producer = producers.shift();
if (!producer) {
throw new Error("No test producer configured");
}
return producer;
}
}

mock.module("kafkajs", () => ({
CompressionTypes: { GZIP: 1 },
Kafka: KafkaMock,
}));

mock.module("./tracing", () => ({ captureError }));

const { disconnectProducer, sendUptimeEvent } = await import("./producer");

const environmentKeys = [
"REDPANDA_BROKER",
"REDPANDA_PASSWORD",
"REDPANDA_SSL",
"REDPANDA_USER",
] as const;
const originalEnvironment = new Map(
environmentKeys.map((key) => [key, process.env[key]])
);

function createProducer(
overrides: Partial<{
connect: () => Promise<void>;
disconnect: () => Promise<void>;
send: () => Promise<void>;
}> = {}
) {
return {
connect: mock(overrides.connect ?? (() => Promise.resolve())),
disconnect: mock(overrides.disconnect ?? (() => Promise.resolve())),
send: mock(overrides.send ?? (() => Promise.resolve())),
};
}

beforeEach(async () => {
await disconnectProducer();
kafkaConfigs.length = 0;
producers.length = 0;
captureError.mockClear();
process.env.REDPANDA_BROKER = "redpanda.test:9092";
delete process.env.REDPANDA_PASSWORD;
delete process.env.REDPANDA_SSL;
delete process.env.REDPANDA_USER;
});

afterAll(async () => {
await disconnectProducer();
for (const key of environmentKeys) {
const value = originalEnvironment.get(key);
if (value === undefined) {
delete process.env[key];
continue;
}
process.env[key] = value;
}
});

describe("sendUptimeEvent", () => {
test("shares one in-flight connection across concurrent cold-start sends", async () => {
let resolveConnection: (() => void) | undefined;
const producer = createProducer({
connect: () =>
new Promise<void>((resolve) => {
resolveConnection = resolve;
}),
});
producers.push(producer);

const sends = Array.from({ length: 20 }, () => sendUptimeEvent({ ok: true }));

expect(producer.connect).toHaveBeenCalledTimes(1);
expect(resolveConnection).toBeDefined();
resolveConnection?.();

await Promise.all(sends);
expect(producer.send).toHaveBeenCalledTimes(20);
expect(producer.send).toHaveBeenCalledWith(
expect.objectContaining({ acks: -1 })
);
expect(kafkaConfigs).toHaveLength(1);
});

test("disconnects a failed producer and reconnects for the next event", async () => {
const failedProducer = createProducer({
send: () => Promise.reject(new Error("broker unavailable")),
});
const recoveredProducer = createProducer();
producers.push(failedProducer, recoveredProducer);

await expect(sendUptimeEvent({ attempt: 1 })).rejects.toThrow(
"broker unavailable"
);
expect(failedProducer.disconnect).toHaveBeenCalledTimes(1);

await expect(sendUptimeEvent({ attempt: 2 })).resolves.toBeUndefined();
expect(recoveredProducer.connect).toHaveBeenCalledTimes(1);
expect(recoveredProducer.send).toHaveBeenCalledTimes(1);
expect(kafkaConfigs).toHaveLength(2);
});

test("reconnects after a failed cold-start connection", async () => {
const failedProducer = createProducer({
connect: () => Promise.reject(new Error("broker unavailable")),
});
const recoveredProducer = createProducer();
producers.push(failedProducer, recoveredProducer);

await expect(sendUptimeEvent({ attempt: 1 })).rejects.toThrow(
"broker unavailable"
);

await expect(sendUptimeEvent({ attempt: 2 })).resolves.toBeUndefined();
expect(recoveredProducer.connect).toHaveBeenCalledTimes(1);
expect(recoveredProducer.send).toHaveBeenCalledTimes(1);
expect(kafkaConfigs).toHaveLength(2);
});

test("rejects when Kafka is not configured", async () => {
delete process.env.REDPANDA_BROKER;

await expect(sendUptimeEvent({ ok: true })).rejects.toThrow(
"REDPANDA_BROKER not set"
);
expect(kafkaConfigs).toEqual([]);
});

test("uses TLS when configured without SASL credentials", async () => {
process.env.REDPANDA_SSL = "true";
const producer = createProducer();
producers.push(producer);

await expect(sendUptimeEvent({ ok: true })).resolves.toBeUndefined();
expect(kafkaConfigs[0]).toEqual(expect.objectContaining({ ssl: true }));
expect(kafkaConfigs[0]).not.toHaveProperty("sasl");
});
});
Loading
Loading