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
54 changes: 54 additions & 0 deletions apps/basket/src/lib/evlog-basket.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown> = {
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<string, unknown> = {
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<string, unknown> = {
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<string, unknown> = {
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<string, unknown> = {
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");
});
});
75 changes: 20 additions & 55 deletions apps/basket/src/lib/evlog-basket.ts
Original file line number Diff line number Diff line change
@@ -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<DrainContext>({
batch: { size: 50, intervalMs: 5000 },
maxBufferSize: 2000,
})(createAxiomDrain({ apiKey: process.env.AXIOM_TOKEN }));
const batchedAxiomDrain = createBatchedAxiomDrain(process.env.AXIOM_TOKEN);

const batchedSuperlogDrain = createBatchedSuperlogDrain();

Expand All @@ -37,41 +34,22 @@ const devFsDrain = useLocalEvlogFiles
? createFsDrain({ dir: devFsLogsDir, pretty: false })
: null;

const DURATION_MS_REGEX = /^([\d.]+)(ms|s)$/;

function normalizeWideEventForAxiom(event: Record<string, unknown>): 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<string, unknown>): 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<string, unknown>
): 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<void> {
Expand All @@ -81,11 +59,6 @@ export async function basketLoggerDrain(ctx: DrainContext): Promise<void> {

normalizeWideEventForAxiom(ctx.event as Record<string, unknown>);

const durationMs = parseDurationMs(ctx.event.duration);
if (durationMs !== undefined) {
ctx.event.duration_ms = durationMs;
}

if (devFsDrain) {
await devFsDrain(ctx);
}
Expand All @@ -95,16 +68,8 @@ export async function basketLoggerDrain(ctx: DrainContext): Promise<void> {
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<void> {
Expand Down
11 changes: 9 additions & 2 deletions apps/basket/src/lib/producer.delivery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof import("evlog")>("evlog");
return { ...actual, log: { ...actual.log, warn: mockLogWarn } };
});

vi.mock("@databuddy/db/clickhouse", () => ({
clickHouse: {},
TABLE_NAMES: {
Expand Down Expand Up @@ -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 () => {
Expand Down Expand Up @@ -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,
Expand Down
8 changes: 6 additions & 2 deletions apps/basket/src/lib/producer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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),
})
)
),
Expand Down
12 changes: 10 additions & 2 deletions apps/basket/src/lib/structured-errors.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<string> = new Set(
Object.values(BASKET_ERROR_SPEC)
.filter((entry) => entry.status >= 400 && entry.status < 500)
.map((entry) => entry.message)
);

declare module "evlog" {
interface RegisteredErrorCatalogs {
Expand Down
1 change: 1 addition & 0 deletions packages/shared/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
},
Expand Down
45 changes: 45 additions & 0 deletions packages/shared/src/evlog-axiom.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown> = {
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<string, unknown> = {
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<string, unknown> = {
error: { status: 503 },
level: "error",
};

normalizeHttpWideEventForAxiom(event);

expect(event.level).toBe("error");
expect(event.client_http_error).toBeUndefined();
});
});
Loading
Loading