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..80b5c4a9a 100644 --- a/apps/basket/src/lib/evlog-basket.ts +++ b/apps/basket/src/lib/evlog-basket.ts @@ -1,21 +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(); @@ -37,41 +34,22 @@ const devFsDrain = useLocalEvlogFiles ? createFsDrain({ dir: devFsLogsDir, pretty: false }) : null; -const DURATION_MS_REGEX = /^([\d.]+)(ms|s)$/; - -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; +function isBasketClientHttpError(event: Record): boolean { + const status = event.http_status; if (typeof status === "number" && status >= 400 && status < 500) { - event.level = "warn"; - event.client_http_error = true; + return true; } + const message = event.error_message; + return typeof message === "string" && CLIENT_ERROR_MESSAGES.has(message); } -function parseDurationMs(duration: unknown): number | undefined { - if (typeof duration !== "string") { - return; - } - const match = duration.match(DURATION_MS_REGEX); - if (!match?.[1]) { - return; +export function normalizeWideEventForAxiom( + event: Record +): void { + normalizeSharedWideEventForAxiom(event); + if (isBasketClientHttpError(event)) { + downgradeClientHttpError(event); } - return match[2] === "s" - ? Math.round(Number.parseFloat(match[1]) * 1000) - : Math.round(Number.parseFloat(match[1])); } export async function basketLoggerDrain(ctx: DrainContext): Promise { @@ -81,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); } @@ -95,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/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), }) ) ), 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 { 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)); +}