From 285fc34149186db6ba9ace9006ac37a97f548209 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 4 Aug 2026 15:08:48 +0300 Subject: [PATCH 1/3] fix(basket): classify expected client errors as warn in axiom drain Expected 4xx conditions (quota exceeded, rate limits, invalid client IDs) were logged at error level because the evlog elysia plugin logs every thrown error as error and drops the status. Derive the client error message set from the error catalog and downgrade those events to warn in the drain, keyed on http_status or the catalog message. --- apps/basket/src/lib/evlog-basket.test.ts | 54 ++++++++++++++++++++++++ apps/basket/src/lib/evlog-basket.ts | 26 ++++++------ apps/basket/src/lib/structured-errors.ts | 12 +++++- 3 files changed, 78 insertions(+), 14 deletions(-) create mode 100644 apps/basket/src/lib/evlog-basket.test.ts diff --git a/apps/basket/src/lib/evlog-basket.test.ts b/apps/basket/src/lib/evlog-basket.test.ts new file mode 100644 index 000000000..a1e6ea62b --- /dev/null +++ b/apps/basket/src/lib/evlog-basket.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, test } from "vitest"; +import { normalizeWideEventForAxiom } from "./evlog-basket"; + +describe("normalizeWideEventForAxiom", () => { + test("downgrades a catalog client error to warn", () => { + const event: Record = { + level: "error", + error_message: "Event quota exceeded", + }; + normalizeWideEventForAxiom(event); + expect(event.level).toBe("warn"); + expect(event.client_http_error).toBe(true); + }); + + test("downgrades based on a 4xx http_status", () => { + const event: Record = { + level: "error", + http_status: 429, + }; + normalizeWideEventForAxiom(event); + expect(event.level).toBe("warn"); + expect(event.client_http_error).toBe(true); + }); + + test("keeps a 5xx catalog error at error level", () => { + const event: Record = { + level: "error", + error_message: "Website lookup temporarily unavailable", + }; + normalizeWideEventForAxiom(event); + expect(event.level).toBe("error"); + expect(event.client_http_error).toBeUndefined(); + }); + + test("keeps an unknown error at error level", () => { + const event: Record = { + level: "error", + error_message: "Failed to get website by ID V2", + }; + normalizeWideEventForAxiom(event); + expect(event.level).toBe("error"); + }); + + test("flattens a string error onto error_message", () => { + const event: Record = { + level: "error", + error: "Event quota exceeded", + }; + normalizeWideEventForAxiom(event); + expect(event.error).toBeUndefined(); + expect(event.error_message).toBe("Event quota exceeded"); + expect(event.level).toBe("warn"); + }); +}); diff --git a/apps/basket/src/lib/evlog-basket.ts b/apps/basket/src/lib/evlog-basket.ts index 4a98f1513..71b0e00b0 100644 --- a/apps/basket/src/lib/evlog-basket.ts +++ b/apps/basket/src/lib/evlog-basket.ts @@ -2,6 +2,7 @@ import { dirname, join } from "node:path"; import { fileURLToPath } from "node:url"; import { readBooleanEnv } from "@databuddy/env/boolean"; import { createBatchedSuperlogDrain } from "@databuddy/shared/evlog-superlog"; +import { CLIENT_ERROR_MESSAGES } from "@lib/structured-errors"; import type { DrainContext, EnrichContext } from "evlog"; import { createAxiomDrain } from "evlog/axiom"; import { @@ -39,23 +40,24 @@ const devFsDrain = useLocalEvlogFiles const DURATION_MS_REGEX = /^([\d.]+)(ms|s)$/; -function normalizeWideEventForAxiom(event: Record): void { +function isClientHttpError(event: Record): boolean { + const status = event.http_status; + if (typeof status === "number" && status >= 400 && status < 500) { + return true; + } + const message = event.error_message; + return typeof message === "string" && CLIENT_ERROR_MESSAGES.has(message); +} + +export function normalizeWideEventForAxiom( + event: Record +): void { if (typeof event.error === "string") { event.error_message = event.error; event.error = undefined; } - if (event.level !== "error") { - return; - } - - const err = event.error; - if (!err || typeof err !== "object" || Array.isArray(err)) { - return; - } - - const status = (err as { status?: number }).status; - if (typeof status === "number" && status >= 400 && status < 500) { + if (event.level === "error" && isClientHttpError(event)) { event.level = "warn"; event.client_http_error = true; } diff --git a/apps/basket/src/lib/structured-errors.ts b/apps/basket/src/lib/structured-errors.ts index 810b5c99b..e73cc4e39 100644 --- a/apps/basket/src/lib/structured-errors.ts +++ b/apps/basket/src/lib/structured-errors.ts @@ -1,7 +1,7 @@ import { createError, defineErrorCatalog, EvlogError, parseError } from "evlog"; import type { z } from "zod"; -export const basketErrorCatalog = defineErrorCatalog("basket", { +const BASKET_ERROR_SPEC = { TRACK_PAYLOAD_TOO_LARGE: { message: "Payload too large", status: 413, @@ -176,7 +176,15 @@ export const basketErrorCatalog = defineErrorCatalog("basket", { why: "The JSON did not match the expected event shape.", fix: "Correct the fields listed in errors and retry.", }, -}); +} as const; + +export const basketErrorCatalog = defineErrorCatalog("basket", BASKET_ERROR_SPEC); + +export const CLIENT_ERROR_MESSAGES: ReadonlySet = new Set( + Object.values(BASKET_ERROR_SPEC) + .filter((entry) => entry.status >= 400 && entry.status < 500) + .map((entry) => entry.message) +); declare module "evlog" { interface RegisteredErrorCatalogs { From 24e77ab6523ecf95912fabf41251f3de90f2825f Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 4 Aug 2026 15:08:58 +0300 Subject: [PATCH 2/3] fix(basket): log redpanda to clickhouse fallback as warn, not error The fallback recovers the delivery via direct ClickHouse insert, so a connection failure that is handled should not surface as an error in Axiom. Log it via log.warn and assert the new channel in the test. --- apps/basket/src/lib/producer.delivery.test.ts | 11 +++++++++-- apps/basket/src/lib/producer.ts | 8 ++++++-- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/apps/basket/src/lib/producer.delivery.test.ts b/apps/basket/src/lib/producer.delivery.test.ts index 845021df6..9e6ac4d12 100644 --- a/apps/basket/src/lib/producer.delivery.test.ts +++ b/apps/basket/src/lib/producer.delivery.test.ts @@ -4,10 +4,16 @@ import type { Admin, Producer } from "kafkajs"; import { beforeEach, describe, expect, test, vi } from "vitest"; import type { ProducerConfig } from "./producer"; -const { mockCaptureError } = vi.hoisted(() => ({ +const { mockCaptureError, mockLogWarn } = vi.hoisted(() => ({ mockCaptureError: vi.fn(), + mockLogWarn: vi.fn(), })); +vi.mock("evlog", async () => { + const actual = await vi.importActual("evlog"); + return { ...actual, log: { ...actual.log, warn: mockLogWarn } }; +}); + vi.mock("@databuddy/db/clickhouse", () => ({ clickHouse: {}, TABLE_NAMES: { @@ -73,6 +79,7 @@ async function makeEffects( describe("producer delivery guarantees", () => { beforeEach(() => { mockCaptureError.mockClear(); + mockLogWarn.mockClear(); }); test("resolves a core send only after direct ClickHouse fallback succeeds", async () => { @@ -745,7 +752,7 @@ describe("producer delivery guarantees", () => { await Promise.all(deliveries); expect(insert).toHaveBeenCalledTimes(50); - expect(mockCaptureError).toHaveBeenCalledTimes(1); + expect(mockLogWarn).toHaveBeenCalledTimes(1); expect(await Effect.runPromise(effects.stats)).toMatchObject({ connected: false, connecting: false, diff --git a/apps/basket/src/lib/producer.ts b/apps/basket/src/lib/producer.ts index 48f133134..c73fd958f 100644 --- a/apps/basket/src/lib/producer.ts +++ b/apps/basket/src/lib/producer.ts @@ -5,7 +5,7 @@ import { readBooleanEnv } from "@databuddy/env/boolean"; import { captureError, record } from "@lib/tracing"; import { PRODUCER_DRAIN_TIMEOUT_MS } from "@lib/shutdown-budget"; import { Data, Deferred, Effect, Layer, ManagedRuntime, Ref } from "effect"; -import { createError } from "evlog"; +import { createError, log } from "evlog"; import { type Admin, CompressionTypes, Kafka, type Producer } from "kafkajs"; function stringifyEvent(event: unknown): string { @@ -569,9 +569,13 @@ function makeProducerEffects( })).pipe( Effect.tap(() => Effect.sync(() => - captureError(err.cause, { + log.warn({ message: "Redpanda connection failed, using ClickHouse fallback", + error_message: + err.cause instanceof Error + ? err.cause.message + : String(err.cause), }) ) ), From 021e11eacf2d135a619f8634fc450c13747a937b Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 4 Aug 2026 15:20:58 +0300 Subject: [PATCH 3/3] refactor(basket): extract axiom drain helpers to @databuddy/shared/evlog-axiom Move the batched Axiom drain, HTTP enrichers, string/duration normalization, and the client-error downgrade into a shared module so basket (and future services) reuse them. Basket keeps its catalog-driven client-error detection and delegates the mechanics to the shared helpers. --- apps/basket/src/lib/evlog-basket.ts | 61 +++----------- packages/shared/package.json | 1 + packages/shared/src/evlog-axiom.test.ts | 45 +++++++++++ packages/shared/src/evlog-axiom.ts | 103 ++++++++++++++++++++++++ 4 files changed, 161 insertions(+), 49 deletions(-) create mode 100644 packages/shared/src/evlog-axiom.test.ts create mode 100644 packages/shared/src/evlog-axiom.ts diff --git a/apps/basket/src/lib/evlog-basket.ts b/apps/basket/src/lib/evlog-basket.ts index 71b0e00b0..80b5c4a9a 100644 --- a/apps/basket/src/lib/evlog-basket.ts +++ b/apps/basket/src/lib/evlog-basket.ts @@ -1,22 +1,18 @@ import { dirname, join } from "node:path"; import { fileURLToPath } from "node:url"; import { readBooleanEnv } from "@databuddy/env/boolean"; +import { + createBatchedAxiomDrain, + downgradeClientHttpError, + enrichHttpWideEvent, + normalizeWideEventForAxiom as normalizeSharedWideEventForAxiom, +} from "@databuddy/shared/evlog-axiom"; import { createBatchedSuperlogDrain } from "@databuddy/shared/evlog-superlog"; import { CLIENT_ERROR_MESSAGES } from "@lib/structured-errors"; import type { DrainContext, EnrichContext } from "evlog"; -import { createAxiomDrain } from "evlog/axiom"; -import { - createRequestSizeEnricher, - createTraceContextEnricher, - createUserAgentEnricher, -} from "evlog/enrichers"; import { createFsDrain } from "evlog/fs"; -import { createDrainPipeline } from "evlog/pipeline"; -const batchedAxiomDrain = createDrainPipeline({ - batch: { size: 50, intervalMs: 5000 }, - maxBufferSize: 2000, -})(createAxiomDrain({ apiKey: process.env.AXIOM_TOKEN })); +const batchedAxiomDrain = createBatchedAxiomDrain(process.env.AXIOM_TOKEN); const batchedSuperlogDrain = createBatchedSuperlogDrain(); @@ -38,9 +34,7 @@ const devFsDrain = useLocalEvlogFiles ? createFsDrain({ dir: devFsLogsDir, pretty: false }) : null; -const DURATION_MS_REGEX = /^([\d.]+)(ms|s)$/; - -function isClientHttpError(event: Record): boolean { +function isBasketClientHttpError(event: Record): boolean { const status = event.http_status; if (typeof status === "number" && status >= 400 && status < 500) { return true; @@ -52,30 +46,12 @@ function isClientHttpError(event: Record): boolean { export function normalizeWideEventForAxiom( event: Record ): void { - if (typeof event.error === "string") { - event.error_message = event.error; - event.error = undefined; - } - - if (event.level === "error" && isClientHttpError(event)) { - event.level = "warn"; - event.client_http_error = true; + normalizeSharedWideEventForAxiom(event); + if (isBasketClientHttpError(event)) { + downgradeClientHttpError(event); } } -function parseDurationMs(duration: unknown): number | undefined { - if (typeof duration !== "string") { - return; - } - const match = duration.match(DURATION_MS_REGEX); - if (!match?.[1]) { - return; - } - return match[2] === "s" - ? Math.round(Number.parseFloat(match[1]) * 1000) - : Math.round(Number.parseFloat(match[1])); -} - export async function basketLoggerDrain(ctx: DrainContext): Promise { if (ctx.event.method === "OPTIONS") { return; @@ -83,11 +59,6 @@ export async function basketLoggerDrain(ctx: DrainContext): Promise { normalizeWideEventForAxiom(ctx.event as Record); - const durationMs = parseDurationMs(ctx.event.duration); - if (durationMs !== undefined) { - ctx.event.duration_ms = durationMs; - } - if (devFsDrain) { await devFsDrain(ctx); } @@ -97,16 +68,8 @@ export async function basketLoggerDrain(ctx: DrainContext): Promise { batchedSuperlogDrain?.(ctx); } -const enrichers = [ - createUserAgentEnricher(), - createRequestSizeEnricher(), - createTraceContextEnricher(), -] as const; - export function enrichBasketWideEvent(ctx: EnrichContext): void { - for (const enricher of enrichers) { - enricher(ctx); - } + enrichHttpWideEvent(ctx); } export async function flushBatchedAxiomDrain(): Promise { diff --git a/packages/shared/package.json b/packages/shared/package.json index 502c31778..4be611517 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -31,6 +31,7 @@ "./stripe-webhooks": "./src/stripe-webhooks.ts", "./audit": "./src/audit.ts", "./evlog-fields": "./src/evlog-fields.ts", + "./evlog-axiom": "./src/evlog-axiom.ts", "./evlog-redaction": "./src/evlog-redaction.ts", "./evlog-superlog": "./src/evlog-superlog.ts" }, diff --git a/packages/shared/src/evlog-axiom.test.ts b/packages/shared/src/evlog-axiom.test.ts new file mode 100644 index 000000000..4b0b4d394 --- /dev/null +++ b/packages/shared/src/evlog-axiom.test.ts @@ -0,0 +1,45 @@ +import { describe, expect, test } from "bun:test"; +import { + normalizeHttpWideEventForAxiom, + normalizeWideEventForAxiom, +} from "./evlog-axiom"; + +describe("normalizeWideEventForAxiom", () => { + test("moves string errors and adds a numeric duration", () => { + const event: Record = { + error: "request failed", + duration: "1.25s", + level: "error", + }; + + normalizeWideEventForAxiom(event); + + expect(event.error).toBeUndefined(); + expect(event.error_message).toBe("request failed"); + expect(event.duration_ms).toBe(1250); + }); + + test("downgrades structured 4xx errors", () => { + const event: Record = { + error: { status: 429 }, + level: "error", + }; + + normalizeHttpWideEventForAxiom(event); + + expect(event.level).toBe("warn"); + expect(event.client_http_error).toBe(true); + }); + + test("keeps structured server errors at error level", () => { + const event: Record = { + error: { status: 503 }, + level: "error", + }; + + normalizeHttpWideEventForAxiom(event); + + expect(event.level).toBe("error"); + expect(event.client_http_error).toBeUndefined(); + }); +}); diff --git a/packages/shared/src/evlog-axiom.ts b/packages/shared/src/evlog-axiom.ts new file mode 100644 index 000000000..1bd811ead --- /dev/null +++ b/packages/shared/src/evlog-axiom.ts @@ -0,0 +1,103 @@ +import type { DrainContext, EnrichContext } from "evlog"; +import { createAxiomDrain } from "evlog/axiom"; +import { + createRequestSizeEnricher, + createTraceContextEnricher, + createUserAgentEnricher, +} from "evlog/enrichers"; +import { createDrainPipeline } from "evlog/pipeline"; + +const createAxiomPipeline = createDrainPipeline({ + batch: { size: 50, intervalMs: 5000 }, + maxBufferSize: 2000, +}); + +const durationPattern = /^([\d.]+)(ms|s)$/; + +const httpEnrichers = [ + createUserAgentEnricher(), + createRequestSizeEnricher(), + createTraceContextEnricher(), +] as const; + +export function createBatchedAxiomDrain( + apiKey: string | undefined, + dataset?: string +) { + return createAxiomPipeline( + createAxiomDrain({ + apiKey, + ...(dataset ? { dataset } : {}), + }) + ); +} + +export function normalizeWideEventForAxiom( + event: Record +): void { + if (typeof event.error === "string") { + event.error_message = event.error; + event.error = undefined; + } + + const durationMs = parseDurationMs(event.duration); + if (durationMs !== undefined) { + event.duration_ms = durationMs; + } +} + +export function normalizeHttpWideEventForAxiom( + event: Record +): void { + normalizeWideEventForAxiom(event); + if (hasClientHttpError(event)) { + downgradeClientHttpError(event); + } +} + +export function downgradeClientHttpError(event: Record): void { + if (event.level !== "error") { + return; + } + + event.level = "warn"; + event.client_http_error = true; +} + +export function enrichHttpWideEvent(ctx: EnrichContext): void { + for (const enrich of httpEnrichers) { + enrich(ctx); + } +} + +function hasClientHttpError(event: Record): boolean { + if (event.level !== "error") { + return false; + } + + const error = event.error; + if (!isRecord(error)) { + return false; + } + + const status = error.status; + return typeof status === "number" && status >= 400 && status < 500; +} + +function parseDurationMs(duration: unknown): number | undefined { + if (typeof duration !== "string") { + return; + } + + const match = duration.match(durationPattern); + if (!match?.[1]) { + return; + } + + const durationMs = Number.parseFloat(match[1]); + return Math.round(match[2] === "s" ? durationMs * 1000 : durationMs); +} + +function isRecord(value: unknown): value is Record { + return Boolean(value && typeof value === "object" && !Array.isArray(value)); +}