diff --git a/.agents/skills/databuddy-internal/SKILL.md b/.agents/skills/databuddy-internal/SKILL.md index ecce69e90..5575bd5fe 100644 --- a/.agents/skills/databuddy-internal/SKILL.md +++ b/.agents/skills/databuddy-internal/SKILL.md @@ -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 diff --git a/apps/uptime/package.json b/apps/uptime/package.json index 0f9a8b53e..9fdcc9b6b 100644 --- a/apps/uptime/package.json +++ b/apps/uptime/package.json @@ -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" } diff --git a/apps/uptime/src/actions.ts b/apps/uptime/src/actions.ts index 1be24b334..a823a6e5b 100644 --- a/apps/uptime/src/actions.ts +++ b/apps/uptime/src/actions.ts @@ -1,3 +1,4 @@ +import { randomUUID } from "node:crypto"; import { connect } from "node:tls"; import { db } from "@databuddy/db"; import { @@ -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, diff --git a/apps/uptime/src/index.ts b/apps/uptime/src/index.ts index 1a30af701..51c54409f 100644 --- a/apps/uptime/src/index.ts +++ b/apps/uptime/src/index.ts @@ -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"), @@ -27,6 +27,12 @@ initLogger({ sampling: {}, }); +let shuttingDown = false; +let shutdownExitCode = 0; +let uptimeWorker: ReturnType | null = null; +let uptimeDeliveryWorker: ReturnType | null = + null; + process.on("unhandledRejection", (reason, _promise) => { captureError(reason, { process: "unhandledRejection" }); log.error({ @@ -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 | 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) => + 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 | null, + deliveryWorker: ReturnType | 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 | 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", diff --git a/apps/uptime/src/lib/producer.test.ts b/apps/uptime/src/lib/producer.test.ts new file mode 100644 index 000000000..3320e5d06 --- /dev/null +++ b/apps/uptime/src/lib/producer.test.ts @@ -0,0 +1,155 @@ +import { afterAll, beforeEach, describe, expect, mock, test } from "bun:test"; + +const kafkaConfigs: unknown[] = []; +const producers: Array> = []; +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; + disconnect: () => Promise; + send: () => Promise; + }> = {} +) { + 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((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"); + }); +}); diff --git a/apps/uptime/src/lib/producer.ts b/apps/uptime/src/lib/producer.ts index 2ba8ca4e6..8a6969a79 100644 --- a/apps/uptime/src/lib/producer.ts +++ b/apps/uptime/src/lib/producer.ts @@ -1,15 +1,8 @@ import { CompressionTypes, Kafka, type Producer } from "kafkajs"; -import { Context, Data, Effect, Layer } from "effect"; import { captureError } from "./tracing"; const TOPIC = "analytics-uptime-checks"; -class KafkaSendError extends Data.TaggedError("KafkaSendError")<{ - cause: unknown; -}> {} - -const KafkaProducer = Context.Service("KafkaProducer"); - const connectProducer = (): Promise => { const broker = process.env.REDPANDA_BROKER; if (!broker) { @@ -21,11 +14,10 @@ const connectProducer = (): Promise => { const kafka = new Kafka({ brokers: [broker], clientId: "uptime-producer", - ...(username && - password && { - sasl: { mechanism: "scram-sha-256", username, password }, - ssl: process.env.REDPANDA_SSL === "true", - }), + ...(username && password + ? { sasl: { mechanism: "scram-sha-256", username, password } } + : {}), + ...(process.env.REDPANDA_SSL === "true" ? { ssl: true } : {}), }); const producer = kafka.producer({ @@ -37,74 +29,44 @@ const connectProducer = (): Promise => { return producer.connect().then(() => producer); }; -const KafkaProducerLive = Layer.effect( - KafkaProducer, - Effect.acquireRelease( - Effect.tryPromise({ - try: connectProducer, - catch: (cause) => { - captureError(cause, { error_step: "kafka_producer_connect" }); - return cause as Error; - }, - }), - (producer) => - Effect.tryPromise({ - try: () => producer.disconnect(), - catch: (cause) => cause, - }).pipe( - Effect.catch((cause) => { - captureError(cause, { - error_step: "kafka_producer_disconnect", - }); - return Effect.void; - }) - ) - ) -); - -const sendEvent = (event: unknown, key?: string) => - Effect.gen(function* () { - const producer = yield* KafkaProducer; - yield* Effect.tryPromise({ - try: () => - producer.send({ - topic: TOPIC, - messages: [ - { - value: JSON.stringify(event, (_k, v) => - v === undefined ? null : v - ), - key, - }, - ], - compression: CompressionTypes.GZIP, - }), - catch: (cause) => new KafkaSendError({ cause }), - }); - }); - -export { KafkaProducer, KafkaProducerLive, KafkaSendError, sendEvent }; - let singletonProducer: Producer | null = null; -let singletonConnected = false; +let singletonConnection: Promise | null = null; -async function ensureProducer(): Promise { - if (singletonConnected && singletonProducer) { - return singletonProducer; +function ensureProducer(): Promise { + if (singletonProducer) { + return Promise.resolve(singletonProducer); + } + if (singletonConnection) { + return singletonConnection; } - if (!process.env.REDPANDA_BROKER) { - return null; + singletonConnection = connectProducer() + .then((producer) => { + singletonProducer = producer; + return producer; + }) + .catch((error) => { + captureError(error, { error_step: "kafka_producer_connect" }); + singletonProducer = null; + throw error; + }) + .finally(() => { + singletonConnection = null; + }); + + return singletonConnection; +} + +async function resetProducer(producer: Producer): Promise { + if (singletonProducer !== producer) { + return; } + singletonProducer = null; try { - singletonProducer = await connectProducer(); - singletonConnected = true; - return singletonProducer; + await producer.disconnect(); } catch (error) { - captureError(error, { error_step: "kafka_producer_connect" }); - singletonConnected = false; - return null; + captureError(error, { error_step: "kafka_producer_disconnect" }); } } @@ -112,17 +74,16 @@ export async function sendUptimeEvent( event: unknown, key?: string ): Promise { - const p = await ensureProducer(); - if (!p) { - return; - } - + const producer = await ensureProducer(); try { - await p.send({ + await producer.send({ topic: TOPIC, + acks: -1, messages: [ { - value: JSON.stringify(event, (_k, v) => (v === undefined ? null : v)), + value: JSON.stringify(event, (_key, value) => + value === undefined ? null : value + ), key, }, ], @@ -130,18 +91,15 @@ export async function sendUptimeEvent( }); } catch (error) { captureError(error, { error_step: "kafka_producer_send" }); + await resetProducer(producer); + throw error; } } export async function disconnectProducer(): Promise { - if (!singletonProducer) { - return; - } - try { - await singletonProducer.disconnect(); - } catch (error) { - captureError(error, { error_step: "kafka_producer_disconnect" }); + const producer = + singletonProducer ?? (await singletonConnection?.catch(() => null)); + if (producer) { + await resetProducer(producer); } - singletonProducer = null; - singletonConnected = false; } diff --git a/apps/uptime/src/sync-schedulers.test.ts b/apps/uptime/src/sync-schedulers.test.ts new file mode 100644 index 000000000..9fa6382cf --- /dev/null +++ b/apps/uptime/src/sync-schedulers.test.ts @@ -0,0 +1,71 @@ +import { afterAll, beforeEach, describe, expect, mock, test } from "bun:test"; +import * as actualDb from "@databuddy/db"; +import * as actualSchema from "@databuddy/db/schema"; +import * as actualRedis from "@databuddy/redis"; +import * as actualEvlog from "evlog"; + +const monitors = [{ granularity: "five_minutes", id: "schedule-1" }]; +const upsertJobScheduler = mock(async () => undefined); +const dbSelect = mock(() => ({ + from: () => ({ where: async () => monitors }), +})); +const logInfo = mock(() => {}); + +mock.module("@databuddy/db", () => ({ + ...actualDb, + db: { select: dbSelect }, + eq: mock(() => undefined), +})); +mock.module("@databuddy/db/schema", () => ({ + ...actualSchema, + uptimeSchedules: { + granularity: "granularity", + id: "id", + isPaused: "isPaused", + }, +})); +mock.module("@databuddy/redis", () => ({ + ...actualRedis, + getUptimeQueue: () => ({ upsertJobScheduler }), + UPTIME_CHECK_JOB_NAME: "uptime-check", + UPTIME_JOB_OPTIONS: { attempts: 1_000_000 }, + uptimeSchedulerId: (scheduleId: string) => `uptime-${scheduleId}`, +})); +mock.module("evlog", () => ({ + ...actualEvlog, + log: { error: mock(() => {}), info: logInfo }, +})); + +const { syncSchedulers } = await import("./sync-schedulers"); + +beforeEach(() => { + dbSelect.mockClear(); + logInfo.mockClear(); + upsertJobScheduler.mockClear(); +}); + +afterAll(() => { + mock.module("@databuddy/db", () => actualDb); + mock.module("@databuddy/db/schema", () => actualSchema); + mock.module("@databuddy/redis", () => actualRedis); + mock.module("evlog", () => actualEvlog); +}); + +describe("syncSchedulers", () => { + test("upserts every active scheduler with current durable job options", async () => { + await syncSchedulers(); + + expect(upsertJobScheduler).toHaveBeenCalledWith( + "uptime-schedule-1", + { pattern: "*/5 * * * *" }, + { + data: { scheduleId: "schedule-1", trigger: "scheduled" }, + name: "uptime-check", + opts: { attempts: 1_000_000 }, + } + ); + expect(logInfo).toHaveBeenCalledWith( + expect.objectContaining({ failed: 0, total: 1, upserted: 1 }) + ); + }); +}); diff --git a/apps/uptime/src/sync-schedulers.ts b/apps/uptime/src/sync-schedulers.ts index da55ad9e2..4747f5a68 100644 --- a/apps/uptime/src/sync-schedulers.ts +++ b/apps/uptime/src/sync-schedulers.ts @@ -31,15 +31,6 @@ const syncMonitor = ( ) => Effect.gen(function* () { const schedulerId = uptimeSchedulerId(monitor.id); - - const existing = yield* Effect.tryPromise({ - try: () => queue.getJobScheduler(schedulerId), - catch: (cause) => cause, - }); - if (existing) { - return "skipped" as const; - } - const pattern = CRON_GRANULARITIES[monitor.granularity]; if (!pattern) { return yield* Effect.fail( @@ -66,8 +57,6 @@ const syncMonitor = ( ), catch: (cause) => cause, }); - - return "created" as const; }); const syncAll = Effect.gen(function* () { @@ -85,17 +74,12 @@ const syncAll = Effect.gen(function* () { catch: (cause) => cause, }); - const created = yield* Ref.make(0); - const skipped = yield* Ref.make(0); + const upserted = yield* Ref.make(0); const failed = yield* Ref.make(0); for (const monitor of monitors) { yield* syncMonitor(monitor, queue).pipe( - Effect.tap((result) => - result === "created" - ? Ref.update(created, (n) => n + 1) - : Ref.update(skipped, (n) => n + 1) - ), + Effect.tap(() => Ref.update(upserted, (n) => n + 1)), Effect.catch((error) => { if (error instanceof UnknownGranularity) { log.error({ @@ -116,17 +100,12 @@ const syncAll = Effect.gen(function* () { ); } - const [c, s, f] = yield* Effect.all([ - Ref.get(created), - Ref.get(skipped), - Ref.get(failed), - ]); + const [u, f] = yield* Effect.all([Ref.get(upserted), Ref.get(failed)]); log.info({ sync: "scheduler", total: monitors.length, - created: c, - skipped: s, + upserted: u, failed: f, }); }); diff --git a/apps/uptime/src/types.ts b/apps/uptime/src/types.ts index f207bb3eb..0b225fa3d 100644 --- a/apps/uptime/src/types.ts +++ b/apps/uptime/src/types.ts @@ -1,3 +1,5 @@ +import { z } from "zod"; + export const MonitorStatus = { DOWN: 0, UP: 1, @@ -5,30 +7,49 @@ export const MonitorStatus = { MAINTENANCE: 3, } as const; -export interface UptimeData { - attempt: number; - check_type: string; - content_hash: string; - env: string; - error: string; - failure_streak: number; - http_code: number; - json_data?: string; - probe_ip: string; - probe_region: string; - redirect_count: number; - response_bytes: number; - retries: number; - site_id: string; - ssl_expiry: number; - ssl_valid: number; - status: number; - timestamp: number; - total_ms: number; - ttfb_ms: number; - url: string; - user_agent: string; -} +export const uptimeDataSchema = z.object({ + attempt: z.number(), + check_type: z.string(), + content_hash: z.string(), + env: z.string(), + error: z.string(), + event_id: z.string(), + failure_streak: z.number(), + http_code: z.number(), + json_data: z.string().optional(), + probe_ip: z.string(), + probe_region: z.string(), + redirect_count: z.number(), + response_bytes: z.number(), + retries: z.number(), + site_id: z.string(), + ssl_expiry: z.number(), + ssl_valid: z.number(), + status: z.number(), + timestamp: z.number(), + total_ms: z.number(), + ttfb_ms: z.number(), + url: z.string(), + user_agent: z.string(), +}); + +const requiredUnknownSchema = z + .unknown() + .refine((value) => value !== undefined, "Required"); + +export const uptimeCheckJobDataSchema = z + .object({ + delivery: z.object({ event: requiredUnknownSchema }).optional(), + scheduleId: z.string(), + trigger: z.enum(["manual", "scheduled"]), + }) + .passthrough(); + +export const uptimeDeliveryJobDataSchema = z.object({ + event: requiredUnknownSchema, +}); + +export type UptimeData = z.infer; export type ScheduleLookupReason = "not_found" | "malformed" | "transient"; diff --git a/apps/uptime/src/uptime-transition-alerts.test.ts b/apps/uptime/src/uptime-transition-alerts.test.ts index f0a921a48..d41a5f1cc 100644 --- a/apps/uptime/src/uptime-transition-alerts.test.ts +++ b/apps/uptime/src/uptime-transition-alerts.test.ts @@ -16,6 +16,7 @@ const baseUptimeData: UptimeData = { attempt: 1, check_type: "http", content_hash: "", + event_id: "uptime-event-1", env: "production", error: "", failure_streak: 0, diff --git a/apps/uptime/src/worker.test.ts b/apps/uptime/src/worker.test.ts index c1c4e2096..24663be67 100644 --- a/apps/uptime/src/worker.test.ts +++ b/apps/uptime/src/worker.test.ts @@ -5,6 +5,7 @@ import { DEFAULT_UPTIME_WORKER_CONCURRENCY, getUptimeWorkerConcurrency, processUptimeCheck, + processUptimeDeliveryJob, processUptimeJob, type UptimeWorkerDeps, } from "./worker"; @@ -18,11 +19,14 @@ const calls = { cacheBust: boolean | undefined; extractHealth: boolean | undefined; }>, + checkpoint: [] as UptimeData[], + delivery: [] as UptimeData[], email: [] as Array<{ schedule: ScheduleData; data: UptimeData }>, loggerFields: [] as Array>, loggerEmitted: [] as Array, + order: [] as string[], reaped: [] as string[], - send: [] as Array<{ data: UptimeData; monitorId: string }>, + send: [] as Array<{ event: unknown; key: string | undefined }>, }; let lookupResult: @@ -55,6 +59,7 @@ function uptimeData(values: Partial = {}): UptimeData { attempt: 1, check_type: "http", content_hash: "hash", + event_id: "uptime-event-1", env: "test", error: "", failure_streak: 0, @@ -104,6 +109,10 @@ function deps(): UptimeWorkerDeps { error: () => {}, } as never; }, + enqueueUptimeDelivery: async (data) => { + calls.delivery.push(data); + calls.order.push("enqueue"); + }, getPreviousMonitorStatus: async () => previousStatus, isHealthExtractionEnabled: (config) => typeof config === "object" && @@ -117,11 +126,12 @@ function deps(): UptimeWorkerDeps { throw new Error("redis reap blew up"); } }, - sendUptimeEvent: async (data, monitorId) => { - calls.send.push({ data, monitorId }); + sendUptimeEvent: async (event, key) => { + calls.send.push({ event, key }); }, fireTransitionAlerts: async (payload) => { calls.email.push(payload); + calls.order.push("alert"); return { transition_kind: null, alarms_fired: 0 }; }, }; @@ -130,9 +140,12 @@ function deps(): UptimeWorkerDeps { beforeEach(() => { calls.captureError = []; calls.check = []; + calls.checkpoint = []; + calls.delivery = []; calls.email = []; calls.loggerFields = []; calls.loggerEmitted = []; + calls.order = []; calls.reaped = []; calls.send = []; lookupResult = { success: true, data: schedule() }; @@ -145,6 +158,25 @@ async function flushMicrotasks(): Promise { await new Promise((resolve) => setImmediate(resolve)); } +type UptimeEventCheckpoint = (data: UptimeData) => Promise; +const noOpCheckpoint: UptimeEventCheckpoint = async () => {}; + +function processUptimeCheckForTest( + scheduleId: string, + trigger: "manual" | "scheduled", + workerDeps: UptimeWorkerDeps = deps(), + jobMeta?: { id?: string; attempt?: number }, + checkpoint: UptimeEventCheckpoint = noOpCheckpoint +) { + return processUptimeCheck( + scheduleId, + trigger, + workerDeps, + jobMeta, + checkpoint + ); +} + describe("getUptimeWorkerConcurrency", () => { it("keeps the high Bun worker default when no override is configured", () => { expect(getUptimeWorkerConcurrency(undefined)).toBe( @@ -186,18 +218,22 @@ describe("processUptimeCheck", () => { { name: "uptime-check", data: { scheduleId: "schedule-1", trigger: "manual" }, + updateData: async (data) => { + calls.checkpoint.push(data.delivery?.event as UptimeData); + }, }, deps() ); expect(calls.check).toHaveLength(1); + expect(calls.checkpoint).toEqual([uptimeData()]); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ uptime_trigger: "manual" }) ); }); it("runs a scheduled check and emits events, status, and transition email work", async () => { - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); expect(calls.check).toEqual([ { @@ -208,10 +244,9 @@ describe("processUptimeCheck", () => { extractHealth: true, }, ]); - expect(calls.send).toEqual([ - { data: uptimeData(), monitorId: "website-1" }, - ]); + expect(calls.delivery).toEqual([uptimeData()]); expect(calls.email).toHaveLength(1); + expect(calls.order).toEqual(["enqueue", "alert"]); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ schedule_id: "schedule-1", @@ -229,6 +264,7 @@ describe("processUptimeCheck", () => { ); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ + event_id: "uptime-event-1", outcome: "up", previous_uptime_status: 0, ttfb_ms: 10, @@ -236,7 +272,7 @@ describe("processUptimeCheck", () => { }) ); expect(calls.loggerFields).toContainEqual( - expect.objectContaining({ kafka_sent: true }) + expect.objectContaining({ delivery_queue_admitted: true }) ); expect(calls.loggerEmitted).toHaveLength(1); }); @@ -244,7 +280,7 @@ describe("processUptimeCheck", () => { it("records -1 when no previous monitor status exists", async () => { previousStatus = undefined; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ previous_uptime_status: -1 }) @@ -257,7 +293,7 @@ describe("processUptimeCheck", () => { data: schedule({ website: null, websiteId: null, timeout: null }), }; - await processUptimeCheck("schedule-only", "manual", deps()); + await processUptimeCheckForTest("schedule-only", "manual", deps()); expect(calls.check).toEqual([ { @@ -285,10 +321,10 @@ describe("processUptimeCheck", () => { it("skips paused schedules without running the check", async () => { lookupResult = { success: true, data: schedule({ isPaused: true }) }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); expect(calls.check).toEqual([]); - expect(calls.send).toEqual([]); + expect(calls.delivery).toEqual([]); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ organization_id: "org-1" }) ); @@ -301,7 +337,7 @@ describe("processUptimeCheck", () => { it("skips missing schedules without throwing", async () => { lookupResult = { success: false, error: "not found" }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); expect(calls.check).toEqual([]); expect(calls.loggerFields).toContainEqual( @@ -320,7 +356,7 @@ describe("processUptimeCheck", () => { reason: "not_found", }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); await flushMicrotasks(); expect(calls.reaped).toEqual(["schedule-1"]); @@ -339,7 +375,7 @@ describe("processUptimeCheck", () => { reason: "malformed", }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); await flushMicrotasks(); expect(calls.reaped).toEqual(["schedule-1"]); @@ -355,7 +391,7 @@ describe("processUptimeCheck", () => { reason: "transient", }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); await flushMicrotasks(); expect(calls.reaped).toEqual([]); @@ -367,7 +403,7 @@ describe("processUptimeCheck", () => { it("does NOT reap when reason is missing on legacy failures (fail-open)", async () => { lookupResult = { success: false, error: "boom" }; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); await flushMicrotasks(); expect(calls.reaped).toEqual([]); @@ -381,7 +417,7 @@ describe("processUptimeCheck", () => { }; reapBehaviour = "throw"; - await processUptimeCheck("schedule-1", "scheduled", deps()); + await processUptimeCheckForTest("schedule-1", "scheduled", deps()); await flushMicrotasks(); expect(calls.reaped).toEqual(["schedule-1"]); @@ -406,7 +442,7 @@ describe("processUptimeCheck", () => { checkResult = { success: false, error: "timeout" }; await expect( - processUptimeCheck("schedule-1", "scheduled", deps()) + processUptimeCheckForTest("schedule-1", "scheduled", deps()) ).rejects.toThrow("timeout"); expect(calls.loggerFields).toContainEqual( expect.objectContaining({ @@ -417,20 +453,161 @@ describe("processUptimeCheck", () => { expect(calls.loggerEmitted).toHaveLength(1); }); - it("captures producer errors on the wide event without failing the job", async () => { + it("persists the exact event before enqueueing it for delivery", async () => { + await processUptimeCheckForTest( + "schedule-1", + "manual", + deps(), + undefined, + async (data) => { + calls.checkpoint.push(data); + calls.order.push("checkpoint"); + } + ); + + expect(calls.checkpoint).toEqual([uptimeData()]); + expect(calls.delivery).toEqual([uptimeData()]); + expect(calls.order).toEqual(["checkpoint", "enqueue", "alert"]); + }); + + it("retries the source job when the durable checkpoint fails", async () => { + await expect( + processUptimeCheckForTest( + "schedule-1", + "manual", + deps(), + undefined, + async () => { + throw new Error("redis unavailable"); + } + ) + ).rejects.toThrow("redis unavailable"); + + expect(calls.delivery).toEqual([]); + expect(calls.email).toEqual([]); + expect(calls.captureError).toContainEqual( + expect.objectContaining({ + context: expect.objectContaining({ + error_step: "uptime_delivery_checkpoint", + event_id: "uptime-event-1", + }), + }) + ); + }); + + it("retries the source job when delivery queue admission fails", async () => { + const failingDeps = deps(); + failingDeps.enqueueUptimeDelivery = async () => { + throw new Error("redis unavailable"); + }; + + await expect( + processUptimeCheckForTest("schedule-1", "manual", failingDeps) + ).rejects.toThrow("redis unavailable"); + + expect(calls.email).toEqual([]); + expect(calls.captureError).toContainEqual( + expect.objectContaining({ + context: expect.objectContaining({ + error_step: "uptime_delivery_enqueue", + event_id: "uptime-event-1", + }), + }) + ); + }); + + it("replays a checkpointed event without running another probe", async () => { + await processUptimeJob( + { + name: "uptime-check", + data: { + delivery: { event: uptimeData() }, + scheduleId: "schedule-1", + trigger: "scheduled", + }, + updateData: async () => {}, + }, + deps() + ); + + expect(calls.check).toEqual([]); + expect(calls.delivery).toEqual([uptimeData()]); + expect(calls.email).toHaveLength(1); + }); + + it("rejects malformed checkpointed delivery payloads before replaying", async () => { + await expect( + processUptimeJob( + { + name: "uptime-check", + data: { + delivery: { event: { ...uptimeData(), http_code: "200" } }, + scheduleId: "schedule-1", + trigger: "scheduled", + }, + updateData: async () => {}, + }, + deps() + ) + ).rejects.toThrow("Invalid persisted uptime delivery payload"); + + expect(calls.check).toEqual([]); + expect(calls.delivery).toEqual([]); + }); + + it("retries a delivery job when Redpanda rejects it", async () => { const failingDeps = deps(); failingDeps.sendUptimeEvent = async () => { - throw new Error("producer unavailable"); + throw new Error("Redpanda send failed"); }; - await processUptimeCheck("schedule-1", "manual", failingDeps); + await expect( + processUptimeDeliveryJob( + { + data: { event: uptimeData() }, + id: "uptime-delivery-uptime-event-1", + name: "uptime-event-delivery", + }, + failingDeps + ) + ).rejects.toThrow("Redpanda send failed"); - expect(calls.loggerFields).toContainEqual( + expect(calls.captureError).toContainEqual( expect.objectContaining({ - kafka_sent: false, - kafka_error: "producer unavailable", + context: expect.objectContaining({ + error_step: "uptime_delivery_send", + event_id: "uptime-event-1", + }), }) ); - expect(calls.loggerEmitted).toHaveLength(1); + }); + + it("delivers the checkpointed payload without its relay-only ID", async () => { + await processUptimeDeliveryJob( + { + data: { event: uptimeData() }, + id: "uptime-delivery-uptime-event-1", + name: "uptime-event-delivery", + }, + deps() + ); + + const { event_id: _eventId, ...event } = uptimeData(); + expect(calls.send).toEqual([{ event, key: "website-1" }]); + }); + + it("rejects malformed delivery payloads before sending", async () => { + await expect( + processUptimeDeliveryJob( + { + data: { event: { ...uptimeData(), site_id: 1 } }, + id: "uptime-delivery-uptime-event-1", + name: "uptime-event-delivery", + }, + deps() + ) + ).rejects.toThrow("Invalid uptime delivery payload"); + + expect(calls.captureError).toEqual([]); }); }); diff --git a/apps/uptime/src/worker.ts b/apps/uptime/src/worker.ts index d8d2a5fab..8a8bb28a1 100644 --- a/apps/uptime/src/worker.ts +++ b/apps/uptime/src/worker.ts @@ -1,13 +1,18 @@ import { getBullMQWorkerConnectionOptions, + getUptimeDeliveryQueue, getUptimeQueue, type UptimeCheckJobData, + type UptimeDeliveryJobData, UPTIME_CHECK_JOB_NAME, + UPTIME_DELIVERY_JOB_NAME, + UPTIME_DELIVERY_QUEUE_NAME, UPTIME_JOB_TIMEOUT_MS, UPTIME_QUEUE_NAME, + uptimeDeliveryJobId, uptimeSchedulerId, } from "@databuddy/redis"; -import { Worker } from "bullmq"; +import { type Job, Worker } from "bullmq"; import type { RequestLogger } from "evlog"; import { createLogger, log } from "evlog"; import { Cause, Data, Effect, Exit } from "effect"; @@ -24,6 +29,9 @@ import { MonitorStatus, type ActionResult, type ScheduleLookupReason, + uptimeCheckJobDataSchema, + uptimeDataSchema, + uptimeDeliveryJobDataSchema, type UptimeData, } from "./types"; import { @@ -54,6 +62,10 @@ class CheckFailed extends Data.TaggedError("CheckFailed")<{ message: string; }> {} +class DeliveryHandoffFailed extends Data.TaggedError("DeliveryHandoffFailed")<{ + message: string; +}> {} + export interface UptimeWorkerDeps { captureError: ( error: unknown, @@ -68,6 +80,7 @@ export interface UptimeWorkerDeps { createLogger: ( fields: Record ) => RequestLogger; + enqueueUptimeDelivery: (data: UptimeData) => Promise; fireTransitionAlerts: (options: { schedule: ScheduleData; data: UptimeData; @@ -80,13 +93,20 @@ export interface UptimeWorkerDeps { isHealthExtractionEnabled: (config: unknown) => boolean; lookupSchedule: (scheduleId: string) => Promise>; reapOrphanScheduler: (scheduleId: string) => Promise; - sendUptimeEvent: (data: UptimeData, monitorId: string) => Promise; + sendUptimeEvent: (event: unknown, key?: string) => Promise; } const uptimeWorkerDeps: UptimeWorkerDeps = { captureError, checkUptime, createLogger: (fields) => createLogger(fields), + enqueueUptimeDelivery: async (data) => { + await getUptimeDeliveryQueue().add( + UPTIME_DELIVERY_JOB_NAME, + { event: data }, + { jobId: uptimeDeliveryJobId(data.event_id) } + ); + }, getPreviousMonitorStatus, isHealthExtractionEnabled, lookupSchedule, @@ -96,6 +116,7 @@ const uptimeWorkerDeps: UptimeWorkerDeps = { }; export const DEFAULT_UPTIME_WORKER_CONCURRENCY = 10_000; +const MAX_STALLED_COUNT = 1_000_000; export function getUptimeWorkerConcurrency( value = process.env.UPTIME_WORKER_CONCURRENCY @@ -112,11 +133,23 @@ export function getUptimeWorkerConcurrency( return parsed; } -export interface UptimeWorkerJob { - attemptsMade?: number; - data: UptimeCheckJobData; - id?: string; - name: string; +export type UptimeWorkerJob = Pick< + Job, + "attemptsMade" | "data" | "id" | "name" | "updateData" +>; + +export type UptimeDeliveryWorkerJob = Pick< + Job, + "attemptsMade" | "data" | "id" | "name" +>; + +type UptimeStorageEvent = Omit; + +function toUptimeStorageEvent({ + event_id: _eventId, + ...event +}: UptimeData): UptimeStorageEvent { + return event; } const timed = ( @@ -171,26 +204,24 @@ const fetchPreviousStatus = (monitorId: string, deps: UptimeWorkerDeps) => Effect.orElseSucceed(() => undefined) ); -const publishEvent = ( +type UptimeEventCheckpoint = (data: UptimeData) => Promise; + +const handoffDelivery = ( data: UptimeData, - monitorId: string, - deps: UptimeWorkerDeps, - log: RequestLogger + handoff: () => Promise, + errorStep: "uptime_delivery_checkpoint" | "uptime_delivery_enqueue", + deps: UptimeWorkerDeps ) => Effect.tryPromise({ - try: () => deps.sendUptimeEvent(data, monitorId), - catch: (cause) => cause, - }).pipe( - Effect.tap(() => Effect.sync(() => log.set({ kafka_sent: true }))), - Effect.catch((error) => - Effect.sync(() => - log.set({ - kafka_sent: false, - kafka_error: error instanceof Error ? error.message : "unknown", - }) - ) - ) - ); + try: handoff, + catch: (cause) => { + deps.captureError(cause, { + error_step: errorStep, + event_id: data.event_id, + }); + return new DeliveryHandoffFailed({ message: String(cause) }); + }, + }); const runTransitionAlerts = ( schedule: ScheduleData, @@ -250,7 +281,8 @@ function reapScheduler( const processCheck = ( scheduleId: string, log: RequestLogger, - deps: UptimeWorkerDeps + deps: UptimeWorkerDeps, + checkpoint: UptimeEventCheckpoint ) => Effect.gen(function* () { const schedule = yield* timed( @@ -320,6 +352,7 @@ const processCheck = ( ); log.set({ + event_id: data.event_id, outcome: data.status === MonitorStatus.UP ? "up" : "down", previous_uptime_status: previousStatus === undefined ? -1 : previousStatus, @@ -337,7 +370,30 @@ const processCheck = ( error_message: data.error || "", }); - yield* timed("kafka", publishEvent(data, monitorId, deps, log), log); + // Persist the completed probe before admission. A source-job retry then + // reuses its event ID and timestamp instead of running a replacement check. + yield* timed( + "delivery_checkpoint", + handoffDelivery( + data, + () => checkpoint(data), + "uptime_delivery_checkpoint", + deps + ), + log + ); + + yield* timed( + "delivery_queue_admission", + handoffDelivery( + data, + () => deps.enqueueUptimeDelivery(data), + "uptime_delivery_enqueue", + deps + ), + log + ); + log.set({ delivery_queue_admitted: true }); yield* timed( "transition_email", @@ -349,8 +405,9 @@ const processCheck = ( export async function processUptimeCheck( scheduleId: string, trigger: UptimeCheckJobData["trigger"], - deps: UptimeWorkerDeps = uptimeWorkerDeps, - jobMeta?: { id?: string; attempt?: number } + deps: UptimeWorkerDeps, + jobMeta: { id?: string; attempt?: number } | undefined, + checkpoint: UptimeEventCheckpoint ) { const startedAt = performance.now(); const log = deps.createLogger({ @@ -360,19 +417,81 @@ export async function processUptimeCheck( ...(jobMeta?.attempt ? { job_attempt: jobMeta.attempt } : {}), }); - const exit = await Effect.runPromiseExit(processCheck(scheduleId, log, deps)); + const exit = await Effect.runPromiseExit( + processCheck(scheduleId, log, deps, checkpoint) + ); log.set({ check_duration_ms: Math.round(performance.now() - startedAt) }); log.emit(); if (Exit.isFailure(exit)) { const error = Cause.squash(exit.cause); - if (error instanceof CheckFailed) { + if ( + error instanceof CheckFailed || + error instanceof DeliveryHandoffFailed + ) { throw new Error(error.message); } } } +async function replayPersistedUptimeDelivery( + job: UptimeWorkerJob, + data: UptimeData, + deps: UptimeWorkerDeps +): Promise { + const startedAt = performance.now(); + const log = deps.createLogger({ + schedule_id: job.data.scheduleId, + uptime_trigger: job.data.trigger, + event_id: data.event_id, + delivery_replay: true, + ...(job.id ? { job_id: job.id } : {}), + ...(job.attemptsMade ? { job_attempt: job.attemptsMade } : {}), + }); + + try { + await Effect.runPromise( + timed( + "delivery_queue_admission", + handoffDelivery( + data, + () => deps.enqueueUptimeDelivery(data), + "uptime_delivery_enqueue", + deps + ), + log + ) + ); + log.set({ delivery_queue_admitted: true }); + + const scheduleExit = await Effect.runPromiseExit( + resolveSchedule(job.data.scheduleId, deps) + ); + if (Exit.isSuccess(scheduleExit)) { + await Effect.runPromise( + timed( + "transition_email", + runTransitionAlerts(scheduleExit.value, data, undefined, deps, log), + log + ) + ); + } else { + const error = Cause.squash(scheduleExit.cause); + log.set({ + transition_alert_skipped: true, + transition_alert_skip_reason: + error instanceof Error ? error.message : String(error), + }); + } + } finally { + log.set({ + delivery_replay_duration_ms: Math.round(performance.now() - startedAt), + }); + log.emit(); + } +} + export async function processUptimeJob( job: UptimeWorkerJob, deps: UptimeWorkerDeps = uptimeWorkerDeps @@ -380,10 +499,72 @@ export async function processUptimeJob( if (job.name !== UPTIME_CHECK_JOB_NAME) { throw new Error(`Unknown uptime job: ${job.name}`); } - await processUptimeCheck(job.data.scheduleId, job.data.trigger, deps, { - id: job.id, - attempt: job.attemptsMade, - }); + + const parsedJobData = uptimeCheckJobDataSchema.safeParse(job.data); + if (!parsedJobData.success) { + throw new Error("Invalid uptime job payload"); + } + + const jobData = parsedJobData.data; + const persistedEvent = jobData.delivery?.event; + if (persistedEvent !== undefined) { + const parsedEvent = uptimeDataSchema.safeParse(persistedEvent); + if (!parsedEvent.success) { + throw new Error("Invalid persisted uptime delivery payload"); + } + await replayPersistedUptimeDelivery(job, parsedEvent.data, deps); + return; + } + + if (typeof job.updateData !== "function") { + throw new Error("Uptime job does not support delivery checkpointing"); + } + + await processUptimeCheck( + jobData.scheduleId, + jobData.trigger, + deps, + { + id: job.id, + attempt: job.attemptsMade, + }, + async (data) => + job.updateData({ + ...jobData, + delivery: { event: data }, + }) + ); +} + +export async function processUptimeDeliveryJob( + job: UptimeDeliveryWorkerJob, + deps: UptimeWorkerDeps = uptimeWorkerDeps +): Promise { + if (job.name !== UPTIME_DELIVERY_JOB_NAME) { + throw new Error(`Unknown uptime delivery job: ${job.name}`); + } + + const parsedJobData = uptimeDeliveryJobDataSchema.safeParse(job.data); + if (!parsedJobData.success) { + throw new Error("Invalid uptime delivery job payload"); + } + + const parsedEvent = uptimeDataSchema.safeParse(parsedJobData.data.event); + if (!parsedEvent.success) { + throw new Error("Invalid uptime delivery payload"); + } + + const data = parsedEvent.data; + try { + await deps.sendUptimeEvent(toUptimeStorageEvent(data), data.site_id); + } catch (error) { + deps.captureError(error, { + error_step: "uptime_delivery_send", + event_id: data.event_id, + job_id: job.id ?? "", + }); + throw error; + } } export function startUptimeWorker() { @@ -394,20 +575,24 @@ export function startUptimeWorker() { connection: getBullMQWorkerConnectionOptions(), concurrency: getUptimeWorkerConcurrency(), lockDuration: UPTIME_JOB_TIMEOUT_MS * 3, + maxStalledCount: MAX_STALLED_COUNT, stalledInterval: UPTIME_JOB_TIMEOUT_MS * 4, } ); worker.on("failed", (job, error) => { const attemptsMade = job?.attemptsMade ?? 0; - const maxAttempts = job?.opts?.attempts ?? 3; + const maxAttempts = job?.opts?.attempts ?? 1_000_000; const isFinalAttempt = attemptsMade >= maxAttempts; + const parsedJobData = job + ? uptimeCheckJobDataSchema.safeParse(job.data) + : undefined; captureError(error, { error_step: "uptime_worker_job_failed", - schedule_id: job?.data.scheduleId ?? "", + schedule_id: parsedJobData?.success ? parsedJobData.data.scheduleId : "", job_id: job?.id ?? "", - trigger: job?.data.trigger ?? "", + trigger: parsedJobData?.success ? parsedJobData.data.trigger : "", attempts_used: attemptsMade, attempts_max: maxAttempts, is_final_attempt: isFinalAttempt, @@ -431,3 +616,51 @@ export function startUptimeWorker() { return worker; } + +export function startUptimeDeliveryWorker() { + const worker = new Worker( + UPTIME_DELIVERY_QUEUE_NAME, + (job) => processUptimeDeliveryJob(job), + { + connection: getBullMQWorkerConnectionOptions(), + concurrency: 1, + lockDuration: UPTIME_JOB_TIMEOUT_MS * 3, + maxStalledCount: MAX_STALLED_COUNT, + stalledInterval: UPTIME_JOB_TIMEOUT_MS * 4, + } + ); + + worker.on("failed", (job, error) => { + const parsedJobData = job + ? uptimeDeliveryJobDataSchema.safeParse(job.data) + : undefined; + const parsedEvent = parsedJobData?.success + ? uptimeDataSchema.safeParse(parsedJobData.data.event) + : undefined; + + captureError(error, { + error_step: "uptime_delivery_worker_job_failed", + event_id: parsedEvent?.success ? parsedEvent.data.event_id : "", + job_id: job?.id ?? "", + attempts_used: job?.attemptsMade ?? 0, + attempts_max: job?.opts?.attempts ?? 1_000_000, + }); + }); + + worker.on("stalled", (jobId) => { + log.warn({ + service: "uptime", + error_step: "uptime_delivery_worker_job_stalled", + error_message: "BullMQ delivery job stalled", + job_id: jobId, + }); + }); + + worker.on("error", (error) => { + captureError(error, { + error_step: "uptime_delivery_worker_error", + }); + }); + + return worker; +} diff --git a/bun.lock b/bun.lock index eb18a8b59..efd4f777b 100644 --- a/bun.lock +++ b/bun.lock @@ -434,6 +434,7 @@ "elysia": "catalog:", "evlog": "catalog:", "kafkajs": "^2.2.4", + "zod": "catalog:", }, }, "apps/video": { @@ -5143,8 +5144,6 @@ "@databuddy/test/drizzle-orm": ["drizzle-orm@1.0.0-rc.1", "", { "peerDependencies": { "@aws-sdk/client-rds-data": ">=3", "@cloudflare/workers-types": ">=4", "@effect/sql-pg": ">=4.0.0-beta.58 || >=4.0.0", "@electric-sql/pglite": ">=0.2.0", "@libsql/client": ">=0.10.0", "@libsql/client-wasm": ">=0.10.0", "@neondatabase/serverless": ">=0.10.0", "@op-engineering/op-sqlite": ">=2", "@opentelemetry/api": "^1.4.1", "@planetscale/database": ">=1.13", "@sinclair/typebox": ">=0.34.8", "@sqlitecloud/drivers": ">=1.0.653", "@tidbcloud/serverless": "*", "@tursodatabase/database": ">=0.2.1", "@tursodatabase/database-common": ">=0.2.1", "@tursodatabase/database-wasm": ">=0.2.1", "@types/better-sqlite3": "*", "@types/mssql": "^9.1.4", "@types/pg": "*", "@types/sql.js": "*", "@upstash/redis": ">=1.34.7", "@vercel/postgres": ">=0.8.0", "@xata.io/client": "*", "arktype": ">=2.0.0", "better-sqlite3": ">=9.3.0", "bun-types": "*", "effect": ">=4.0.0-beta.58 || >=4.0.0", "expo-sqlite": ">=14.0.0", "mssql": "^11.0.1", "mysql2": ">=2", "pg": ">=8", "postgres": ">=3", "sql.js": ">=1", "sqlite3": ">=5", "typebox": ">=1.0.0", "valibot": ">=1.0.0-beta.7", "zod": "^3.25.0 || ^4.0.0" }, "optionalPeers": ["@aws-sdk/client-rds-data", "@cloudflare/workers-types", "@effect/sql-pg", "@electric-sql/pglite", "@libsql/client", "@libsql/client-wasm", "@neondatabase/serverless", "@op-engineering/op-sqlite", "@opentelemetry/api", "@planetscale/database", "@sinclair/typebox", "@sqlitecloud/drivers", "@tidbcloud/serverless", "@tursodatabase/database", "@tursodatabase/database-common", "@tursodatabase/database-wasm", "@types/better-sqlite3", "@types/mssql", "@types/pg", "@types/sql.js", "@upstash/redis", "@vercel/postgres", "@xata.io/client", "arktype", "better-sqlite3", "bun-types", "effect", "expo-sqlite", "mssql", "mysql2", "pg", "postgres", "sql.js", "sqlite3", "typebox", "valibot", "zod"] }, "sha512-jGCqAgxpz+OSHP2jQGooUHBxnFMTYl0TTRSfULBl52VNf7CtyNRnazUi+VdbSxvJrDP2lnIsmUh5O+HhKeSJCg=="], - "@databuddy/uptime/effect": ["effect@4.0.0-beta.66", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.6.0", "find-my-way-ts": "^0.1.6", "ini": "^6.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^1.11.9", "multipasta": "^0.2.7", "toml": "^4.1.1", "uuid": "^13.0.0", "yaml": "^2.8.3" } }, "sha512-4arEr62cziFa8BBVDUwJCJJmaVepXf/kRg7KtC0h8+bufngscrHbwWFhr9c+HonwOF+31U3iD3xUJmw9KzX7Dw=="], - "@dxup/nuxt/tinyglobby": ["tinyglobby@0.2.16", "", { "dependencies": { "fdir": "^6.5.0", "picomatch": "^4.0.4" } }, "sha512-pn99VhoACYR8nFHhxqix+uvsbXineAasWm5ojXoN8xEwK5Kd3/TrhNn1wByuD52UxWRLy8pu+kRMniEi6Eq9Zg=="], "@esbuild-kit/core-utils/esbuild": ["esbuild@0.18.20", "", { "optionalDependencies": { "@esbuild/android-arm": "0.18.20", "@esbuild/android-arm64": "0.18.20", "@esbuild/android-x64": "0.18.20", "@esbuild/darwin-arm64": "0.18.20", "@esbuild/darwin-x64": "0.18.20", "@esbuild/freebsd-arm64": "0.18.20", "@esbuild/freebsd-x64": "0.18.20", "@esbuild/linux-arm": "0.18.20", "@esbuild/linux-arm64": "0.18.20", "@esbuild/linux-ia32": "0.18.20", "@esbuild/linux-loong64": "0.18.20", "@esbuild/linux-mips64el": "0.18.20", "@esbuild/linux-ppc64": "0.18.20", "@esbuild/linux-riscv64": "0.18.20", "@esbuild/linux-s390x": "0.18.20", "@esbuild/linux-x64": "0.18.20", "@esbuild/netbsd-x64": "0.18.20", "@esbuild/openbsd-x64": "0.18.20", "@esbuild/sunos-x64": "0.18.20", "@esbuild/win32-arm64": "0.18.20", "@esbuild/win32-ia32": "0.18.20", "@esbuild/win32-x64": "0.18.20" }, "bin": { "esbuild": "bin/esbuild" } }, "sha512-ceqxoedUrcayh7Y7ZX6NdbbDzGROiyVBgC4PriJThBKSVPWnnFHZAkfI1lJT8QFkOwH4qOS2SJkS4wvpGl8BpA=="], @@ -5977,12 +5976,6 @@ "@databuddy/docs/resend/svix": ["svix@1.92.2", "", { "dependencies": { "standardwebhooks": "1.0.0" } }, "sha512-ZmuA3UVvlnF9EgxlzmPtF7CKjQb64Z6OFlyfdDfU0sdcC7dJa+3aOYX5B9mA+RS6ch1AxBa4UP/l6KmqfGtWBQ=="], - "@databuddy/uptime/effect/ini": ["ini@6.0.0", "", {}, "sha512-IBTdIkzZNOpqm7q3dRqJvMaldXjDHWkEDfrwGEQTs5eaQMWV+djAhR+wahyNNMAa+qpbDUhBMVt4ZKNwpPm7xQ=="], - - "@databuddy/uptime/effect/msgpackr": ["msgpackr@1.11.12", "", { "optionalDependencies": { "msgpackr-extract": "^3.0.2" } }, "sha512-RBdJ1Un7yGlXWajrkxcSa93nvQ0w4zBf60c0yYv7YtBelP8H2FA7XsfBbMHtXKXUMUxH7zV3Zuozh+kUQWhHvg=="], - - "@databuddy/uptime/effect/uuid": ["uuid@13.0.2", "", { "bin": { "uuid": "dist-node/bin/uuid" } }, "sha512-vzi9uRZ926x4XV73S/4qQaTwPXM2JBj6/6lI/byHH1jOpCzb0zDbfytgA9LcN/hzb2l7WQSQnxITOVx5un/wGw=="], - "@esbuild-kit/core-utils/esbuild/@esbuild/android-arm": ["@esbuild/android-arm@0.18.20", "", { "os": "android", "cpu": "arm" }, "sha512-fyi7TDI/ijKKNZTUJAQqiG5T7YjJXgnzkURqmGj13C6dCqckZBLdl4h7bkhHt/t0WP+zO9/zwroDvANaOqO5Sw=="], "@esbuild-kit/core-utils/esbuild/@esbuild/android-arm64": ["@esbuild/android-arm64@0.18.20", "", { "os": "android", "cpu": "arm64" }, "sha512-Nz4rJcchGDtENV0eMKUNa6L12zz2zBDXuhj/Vjh18zGqB44Bi7MBMSXjgunJgjRhCmKOjnPuZp4Mb6OKqtMHLQ=="], @@ -6879,8 +6872,6 @@ "zip-stream/readable-stream/string_decoder": ["string_decoder@1.3.0", "", { "dependencies": { "safe-buffer": "~5.2.0" } }, "sha512-hkRX8U1WjJFd8LsDJ2yQ/wWWxaopEsABU1XfkM8A+j0+85JAGppt16cr1Whg6KIbb4okU6Mql6BOj+uup/wKeA=="], - "@databuddy/uptime/effect/msgpackr/msgpackr-extract": ["msgpackr-extract@3.0.3", "", { "dependencies": { "node-gyp-build-optional-packages": "5.2.2" }, "optionalDependencies": { "@msgpackr-extract/msgpackr-extract-darwin-arm64": "3.0.3", "@msgpackr-extract/msgpackr-extract-darwin-x64": "3.0.3", "@msgpackr-extract/msgpackr-extract-linux-arm": "3.0.3", "@msgpackr-extract/msgpackr-extract-linux-arm64": "3.0.3", "@msgpackr-extract/msgpackr-extract-linux-x64": "3.0.3", "@msgpackr-extract/msgpackr-extract-win32-x64": "3.0.3" }, "bin": { "download-msgpackr-prebuilds": "bin/download-prebuilds.js" } }, "sha512-P0efT1C9jIdVRefqjzOQ9Xml57zpOXnIuS+csaB4MdZbTdmGDLo8XhzBG1N7aO11gKDDkJvBLULeFTo46wwreA=="], - "@inquirer/core/wrap-ansi/string-width/is-fullwidth-code-point": ["is-fullwidth-code-point@3.0.0", "", {}, "sha512-zymm5+u+sCsSWyD9qNaejV3DFvhCKclKdizYaJUuHA83RLjb7nSuGnddCHGv0hk+KY7BMAlsWeK4Ueg6EV6XQg=="], "@mapbox/node-pre-gyp/node-fetch/whatwg-url/tr46": ["tr46@0.0.3", "", {}, "sha512-N3WMsuqV66lT30CrXNbEjx4GEwlow3v6rr4mCcv6prnfwhS01rkgyFdjPNBYd9br7LpXV1+Emh01fHnq2Gdgrw=="], @@ -7171,18 +7162,6 @@ "zip-stream/readable-stream/string_decoder/safe-buffer": ["safe-buffer@5.2.1", "", {}, "sha512-rp3So07KcdmmKbGvgaNxQSJr7bGVSVk5S9Eq1F+ppbRo70+YeaDxkw5Dd8NPN+GD6bjnYm2VuPuCXmpuYvmCXQ=="], - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-darwin-arm64": ["@msgpackr-extract/msgpackr-extract-darwin-arm64@3.0.3", "", { "os": "darwin", "cpu": "arm64" }, "sha512-QZHtlVgbAdy2zAqNA9Gu1UpIuI8Xvsd1v8ic6B2pZmeFnFcMWiPLfWXh7TVw4eGEZ/C9TH281KwhVoeQUKbyjw=="], - - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-darwin-x64": ["@msgpackr-extract/msgpackr-extract-darwin-x64@3.0.3", "", { "os": "darwin", "cpu": "x64" }, "sha512-mdzd3AVzYKuUmiWOQ8GNhl64/IoFGol569zNRdkLReh6LRLHOXxU4U8eq0JwaD8iFHdVGqSy4IjFL4reoWCDFw=="], - - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-linux-arm": ["@msgpackr-extract/msgpackr-extract-linux-arm@3.0.3", "", { "os": "linux", "cpu": "arm" }, "sha512-fg0uy/dG/nZEXfYilKoRe7yALaNmHoYeIoJuJ7KJ+YyU2bvY8vPv27f7UKhGRpY6euFYqEVhxCFZgAUNQBM3nw=="], - - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-linux-arm64": ["@msgpackr-extract/msgpackr-extract-linux-arm64@3.0.3", "", { "os": "linux", "cpu": "arm64" }, "sha512-YxQL+ax0XqBJDZiKimS2XQaf+2wDGVa1enVRGzEvLLVFeqa5kx2bWbtcSXgsxjQB7nRqqIGFIcLteF/sHeVtQg=="], - - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-linux-x64": ["@msgpackr-extract/msgpackr-extract-linux-x64@3.0.3", "", { "os": "linux", "cpu": "x64" }, "sha512-cvwNfbP07pKUfq1uH+S6KJ7dT9K8WOE4ZiAcsrSes+UY55E/0jLYc+vq+DO7jlmqRb5zAggExKm0H7O/CBaesg=="], - - "@databuddy/uptime/effect/msgpackr/msgpackr-extract/@msgpackr-extract/msgpackr-extract-win32-x64": ["@msgpackr-extract/msgpackr-extract-win32-x64@3.0.3", "", { "os": "win32", "cpu": "x64" }, "sha512-x0fWaQtYp4E6sktbsdAqnehxDgEc/VwM7uLsRCYWaiGu0ykYdZPiS8zCWdnjHwyiumousxfBm4SO31eXqwEZhQ=="], - "archiver-utils/glob/minimatch/brace-expansion/balanced-match": ["balanced-match@1.0.2", "", {}, "sha512-3oSeUO0TMV67hN1AmbXsK4yaqU7tjiHlbxRDZOpH0KW9+CeX4bRAaX0Anxt0tx2MrpRpWwQaPwIlISEJhYU5Pw=="], "nuxt/oxc-parser/@oxc-parser/binding-wasm32-wasi/@emnapi/core/@emnapi/wasi-threads": ["@emnapi/wasi-threads@1.2.1", "", { "dependencies": { "tslib": "^2.4.0" } }, "sha512-uTII7OYF+/Mes/MrcIOYp5yOtSMLBWSIoLPpcgwipoiKbli6k322tcoFsxoIIxPDqW01SQGAgko4EzZi2BNv2w=="], diff --git a/packages/ai/src/ai/agents/cache.test.ts b/packages/ai/src/ai/agents/cache.test.ts index 03fe4864c..9c5552f94 100644 --- a/packages/ai/src/ai/agents/cache.test.ts +++ b/packages/ai/src/ai/agents/cache.test.ts @@ -120,6 +120,7 @@ vi.mock("@databuddy/redis", () => ({ getRedisCache: () => mockRedisClient, getUptimeQueue: vi.fn(() => ({})), // Bun keeps this module mock active across later AI tests that import links RPC. + // Keep the full link-cache surface local so those tests never reach Redis. abandonCachedLinkMutation: vi.fn(async () => true), beginCachedLinkMutation: vi.fn(async () => ({ state: "acquired", @@ -171,6 +172,7 @@ vi.mock("@databuddy/redis", () => ({ setCachedLink: vi.fn(async () => undefined), setCachedLinkIfAbsent: vi.fn(async () => true), setCachedLinkNotFound: vi.fn(async () => undefined), + setCachedLinkNotFoundIfAbsent: vi.fn(async () => true), shouldRecordClick: vi.fn(async () => true), shutdownRedis: vi.fn(async () => undefined), streamBufferKey: (id: string) => `stream:${id}`, diff --git a/packages/ai/src/ai/mcp/conversation-store.test.ts b/packages/ai/src/ai/mcp/conversation-store.test.ts index db6f8614f..8c0a15cbf 100644 --- a/packages/ai/src/ai/mcp/conversation-store.test.ts +++ b/packages/ai/src/ai/mcp/conversation-store.test.ts @@ -34,7 +34,12 @@ vi.mock("@databuddy/redis", () => ({ INSIGHTS_QUEUE_ENV_PREFIX: "INSIGHTS", INSIGHTS_QUEUE_NAME: "insights-generation", activeStreamKey: (id: string) => `active:${id}`, + abandonCachedLinkMutation: vi.fn(async () => true), appendStreamChunk: vi.fn(async () => undefined), + beginCachedLinkMutation: vi.fn(async () => ({ + state: "acquired", + token: "test-token", + })), cacheNamespaces: { agentTelemetryWebsiteExists: "agent-telemetry:website-exists", apiKeyByHash: "api-key-by-hash", @@ -95,6 +100,7 @@ vi.mock("@databuddy/redis", () => ({ }, getInsightsQueue: vi.fn(() => ({})), getUptimeQueue: vi.fn(() => ({})), + finishCachedLinkMutation: vi.fn(async () => true), invalidateAgentContextSnapshot: vi.fn(async () => 0), invalidateAgentContextSnapshotsForOwner: vi.fn(async () => 0), invalidateAgentContextSnapshotsForWebsite: vi.fn(async () => 0), @@ -138,7 +144,9 @@ vi.mock("@databuddy/redis", () => ({ redis: mockRedisClient, setActiveStream: vi.fn(async () => undefined), setCachedLink: vi.fn(async () => undefined), + setCachedLinkIfAbsent: vi.fn(async () => true), setCachedLinkNotFound: vi.fn(async () => undefined), + setCachedLinkNotFoundIfAbsent: vi.fn(async () => true), shouldRecordClick: vi.fn(async () => true), shutdownRedis: vi.fn(async () => undefined), streamBufferKey: (id: string) => `stream:${id}`, diff --git a/packages/redis/uptime-queue.test.ts b/packages/redis/uptime-queue.test.ts index ef95845a3..b2bba323f 100644 --- a/packages/redis/uptime-queue.test.ts +++ b/packages/redis/uptime-queue.test.ts @@ -27,10 +27,15 @@ process.env.BULLMQ_REDIS_URL = "redis://queue-user:queue-pass@queue.test:6381/4" const { closeUptimeQueue, + getUptimeDeliveryQueue, getUptimeQueue, UPTIME_CHECK_JOB_NAME, + UPTIME_DELIVERY_JOB_NAME, + UPTIME_DELIVERY_JOB_OPTIONS, + UPTIME_DELIVERY_QUEUE_NAME, UPTIME_JOB_OPTIONS, UPTIME_QUEUE_NAME, + uptimeDeliveryJobId, uptimeImmediateJobId, uptimeSchedulerId, } = await import("./uptime-queue"); @@ -71,15 +76,39 @@ describe("uptime queue", () => { expect(constructorCalls).toHaveLength(1); }); + it("constructs a separate durable delivery queue", () => { + const queue = getUptimeDeliveryQueue(); + + expect(queue).toBeInstanceOf(MockQueue); + expect(constructorCalls).toHaveLength(1); + expect(constructorCalls[0]).toEqual({ + name: UPTIME_DELIVERY_QUEUE_NAME, + options: { + connection: { + host: "queue.test", + port: 6381, + username: "queue-user", + password: "queue-pass", + db: 4, + maxRetriesPerRequest: 1, + }, + defaultJobOptions: UPTIME_DELIVERY_JOB_OPTIONS, + }, + }); + }); + it("closes and resets the singleton", async () => { const first = getUptimeQueue(); + const firstDelivery = getUptimeDeliveryQueue(); await closeUptimeQueue(); const second = getUptimeQueue(); + const secondDelivery = getUptimeDeliveryQueue(); - expect(closeCalls).toBe(1); + expect(closeCalls).toBe(2); expect(second).not.toBe(first); - expect(constructorCalls).toHaveLength(2); + expect(secondDelivery).not.toBe(firstDelivery); + expect(constructorCalls).toHaveLength(4); }); it("uses stable queue constants and namespaced ids", () => { @@ -88,9 +117,25 @@ describe("uptime queue", () => { expect(UPTIME_CHECK_JOB_NAME).toBe("uptime-check"); expect(UPTIME_QUEUE_NAME).toBe("uptime-checks"); + expect(UPTIME_DELIVERY_JOB_NAME).toBe("uptime-event-delivery"); + expect(UPTIME_DELIVERY_QUEUE_NAME).toBe("uptime-event-delivery"); expect(uptimeSchedulerId("schedule-1")).toBe("uptime-schedule-1"); + expect(uptimeDeliveryJobId("event-1")).toBe("uptime-delivery-event-1"); expect(first.startsWith("uptime-manual-schedule-1-")).toBe(true); expect(second.startsWith("uptime-manual-schedule-1-")).toBe(true); expect(first).not.toBe(second); }); + + it("keeps source and delivery jobs retryable through infrastructure outages", () => { + expect(UPTIME_JOB_OPTIONS).toMatchObject({ + attempts: 1_000_000, + backoff: { delay: 30_000, type: "fixed" }, + removeOnFail: false, + }); + expect(UPTIME_DELIVERY_JOB_OPTIONS).toMatchObject({ + attempts: 1_000_000, + backoff: { delay: 30_000, type: "fixed" }, + removeOnFail: false, + }); + }); }); diff --git a/packages/redis/uptime-queue.ts b/packages/redis/uptime-queue.ts index e398cd94a..6032d9117 100644 --- a/packages/redis/uptime-queue.ts +++ b/packages/redis/uptime-queue.ts @@ -3,31 +3,55 @@ import { getBullMQConnectionOptions } from "./bullmq"; export const UPTIME_QUEUE_NAME = "uptime-checks"; export const UPTIME_CHECK_JOB_NAME = "uptime-check"; +export const UPTIME_DELIVERY_QUEUE_NAME = "uptime-event-delivery"; +export const UPTIME_DELIVERY_JOB_NAME = UPTIME_DELIVERY_QUEUE_NAME; export const UPTIME_JOB_TIMEOUT_MS = 30_000; - -export const UPTIME_JOB_OPTIONS = { - attempts: 3, +const UPTIME_RETRY_ATTEMPTS = 1_000_000; +const UPTIME_RETRY_DELAY_MS = 30_000; +const UPTIME_RETRY_OPTIONS = { + attempts: UPTIME_RETRY_ATTEMPTS, backoff: { - type: "exponential", - delay: 2000, + type: "fixed", + delay: UPTIME_RETRY_DELAY_MS, }, + removeOnFail: false, +}; + +/** + * A target that is down produces an UptimeData result and completes normally. + * These retries only cover worker and durable-handoff failures. + */ +export const UPTIME_JOB_OPTIONS = { + ...UPTIME_RETRY_OPTIONS, + // A failed source job can contain the only durable copy of a completed probe. removeOnComplete: { age: 24 * 3600, count: 1000, }, - removeOnFail: { +}; + +export const UPTIME_DELIVERY_JOB_OPTIONS = { + ...UPTIME_RETRY_OPTIONS, + // Keep completed IDs long enough for an ambiguous queue add to stay idempotent. + removeOnComplete: { age: 7 * 24 * 3600, - count: 5000, + count: 10_000, }, }; export interface UptimeCheckJobData { + delivery?: { event: unknown }; scheduleId: string; trigger: "manual" | "scheduled"; } +export interface UptimeDeliveryJobData { + event: unknown; +} + let uptimeQueue: Queue | null = null; +let uptimeDeliveryQueue: Queue | null = null; export function getUptimeQueue(): Queue { uptimeQueue ??= new Queue(UPTIME_QUEUE_NAME, { @@ -38,13 +62,23 @@ export function getUptimeQueue(): Queue { return uptimeQueue; } +export function getUptimeDeliveryQueue(): Queue { + uptimeDeliveryQueue ??= new Queue( + UPTIME_DELIVERY_QUEUE_NAME, + { + connection: getBullMQConnectionOptions(), + defaultJobOptions: UPTIME_DELIVERY_JOB_OPTIONS, + } + ); + + return uptimeDeliveryQueue; +} + export async function closeUptimeQueue(): Promise { - if (!uptimeQueue) { - return; - } - const queue = uptimeQueue; + const queues = [uptimeQueue, uptimeDeliveryQueue]; uptimeQueue = null; - await queue.close(); + uptimeDeliveryQueue = null; + await Promise.all(queues.map((queue) => queue?.close())); } export function uptimeSchedulerId(scheduleId: string): string { @@ -54,3 +88,7 @@ export function uptimeSchedulerId(scheduleId: string): string { export function uptimeImmediateJobId(scheduleId: string): string { return `uptime-manual-${scheduleId}-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`; } + +export function uptimeDeliveryJobId(eventId: string): string { + return `uptime-delivery-${eventId}`; +} diff --git a/packages/rpc/src/services/uptime-scheduler.test.ts b/packages/rpc/src/services/uptime-scheduler.test.ts index f22416ea3..2f4e7a2a4 100644 --- a/packages/rpc/src/services/uptime-scheduler.test.ts +++ b/packages/rpc/src/services/uptime-scheduler.test.ts @@ -2,19 +2,16 @@ import { beforeEach, describe, expect, it, mock } from "bun:test"; const UPTIME_CHECK_JOB_NAME = "uptime-check"; const UPTIME_JOB_OPTIONS = { - attempts: 3, + attempts: 1_000_000, backoff: { - type: "exponential", - delay: 2000, + type: "fixed", + delay: 30_000, }, removeOnComplete: { age: 24 * 3600, count: 1000, }, - removeOnFail: { - age: 7 * 24 * 3600, - count: 5000, - }, + removeOnFail: false, }; const calls: { diff --git a/packages/tracker/src/plugins/pixel.ts b/packages/tracker/src/plugins/pixel.ts index 57a77f616..c5e13f704 100644 --- a/packages/tracker/src/plugins/pixel.ts +++ b/packages/tracker/src/plugins/pixel.ts @@ -2,7 +2,6 @@ import type { BaseTracker } from "../core/tracker"; import type { HttpResult } from "../core/client"; const PIXEL_PATH = "/px.jpg"; -const PIXEL_RETRY_AFTER_MS = 5000; const PIXEL_TYPE_BY_ENDPOINT: Record = { "/": "track", @@ -12,9 +11,12 @@ const PIXEL_TYPE_BY_ENDPOINT: Record = { "/vitals": "web_vitals", "/errors": "error", }; +const MAX_PIXEL_RETRY_DELAY_MS = 30_000; +const PIXEL_LOAD_TIMEOUT_MS = 10_000; -function wait(ms: number): Promise { - return new Promise((resolve) => setTimeout(resolve, ms)); +interface PixelDeliveryResult { + attempts: number; + success: boolean; } function safeStringify(value: unknown): string { @@ -54,21 +56,61 @@ function flattenIntoParams( export function initPixelTracking(tracker: BaseTracker) { tracker.options.enableBatching = false; + let deliveryGeneration = 0; + const activeDeliveryCancels = new Set<() => void>(); + const cancelPendingRequests = tracker.api.cancelPendingRequests.bind( + tracker.api + ); + tracker.api.cancelPendingRequests = () => { + deliveryGeneration += 1; + for (const cancel of [...activeDeliveryCancels]) { + cancel(); + } + cancelPendingRequests(); + }; + + const configuredMaxRetries = tracker.options.maxRetries; const maxRetries = tracker.options.enableRetries === false ? 0 - : (tracker.options.maxRetries ?? 3); - const retryDelay = Math.max( - 0, - tracker.options.initialRetryDelay ?? PIXEL_RETRY_AFTER_MS - ); + : typeof configuredMaxRetries === "number" && + Number.isFinite(configuredMaxRetries) && + configuredMaxRetries >= 0 + ? Math.floor(configuredMaxRetries) + : 3; + const configuredInitialRetryDelay = tracker.options.initialRetryDelay; + const initialRetryDelay = + typeof configuredInitialRetryDelay === "number" && + Number.isFinite(configuredInitialRetryDelay) && + configuredInitialRetryDelay >= 0 + ? configuredInitialRetryDelay + : 500; + + const waitForRetry = (retry: number, generation: number): Promise => + new Promise((resolve) => { + let settled = false; + const finish = (shouldRetry: boolean) => { + if (settled) { + return; + } + settled = true; + clearTimeout(timeout); + activeDeliveryCancels.delete(cancel); + resolve(shouldRetry); + }; + const cancel = () => finish(false); + const timeout = setTimeout( + () => finish(generation === deliveryGeneration), + Math.min(initialRetryDelay * 2 ** retry, MAX_PIXEL_RETRY_DELAY_MS) + ); + activeDeliveryCancels.add(cancel); + }); - const sendOnePixel = async ( + const sendOnePixel = ( eventType: string, - data: Record, - retryCount = 0 - ): Promise<{ attempts: number; success: boolean }> => { + data: Record + ): Promise => { const params = new URLSearchParams(); flattenIntoParams(params, data); @@ -91,26 +133,75 @@ export function initPixelTracking(tracker: BaseTracker) { url.searchParams.append(key, value); }); - const result = await new Promise<{ success: boolean }>((resolve) => { - const img = new Image(); - img.onload = () => resolve({ success: true }); - img.onerror = () => resolve({ success: false }); - img.src = url.toString(); - }); - if (result.success || retryCount >= maxRetries) { - return { attempts: retryCount + 1, success: result.success }; - } - await wait(retryDelay); - return sendOnePixel(eventType, data, retryCount + 1); + const generation = deliveryGeneration; + const load = (): Promise => + new Promise((resolve) => { + let img: HTMLImageElement; + try { + img = new Image(); + } catch { + resolve(false); + return; + } + + let settled = false; + const finish = (success: boolean) => { + if (settled) { + return; + } + settled = true; + clearTimeout(timeout); + activeDeliveryCancels.delete(cancel); + img.onload = null; + img.onerror = null; + resolve(success); + }; + const cancel = () => { + if (settled) { + return; + } + img.onload = null; + img.onerror = null; + img.src = ""; + finish(false); + }; + const timeout = setTimeout(cancel, PIXEL_LOAD_TIMEOUT_MS); + + activeDeliveryCancels.add(cancel); + img.onload = () => finish(true); + img.onerror = () => finish(false); + img.src = url.toString(); + }); + + return (async () => { + for (let retry = 0; retry <= maxRetries; retry += 1) { + if (generation !== deliveryGeneration) { + return { attempts: retry, success: false }; + } + if (await load()) { + return { attempts: retry + 1, success: true }; + } + if (retry < maxRetries && generation === deliveryGeneration) { + const shouldRetry = await waitForRetry(retry, generation); + if (!shouldRetry) { + return { attempts: retry + 1, success: false }; + } + } else if (retry < maxRetries) { + return { attempts: retry + 1, success: false }; + } + } + + return { attempts: maxRetries + 1, success: false }; + })(); }; const sendToPixel = async ( endpoint: string, data: unknown - ): Promise<{ attempts: number; success: boolean }> => { + ): Promise => { const eventType = PIXEL_TYPE_BY_ENDPOINT[endpoint]; if (!eventType) { - return { attempts: 0, success: false }; + return { attempts: 1, success: false }; } if (Array.isArray(data)) { @@ -118,17 +209,17 @@ export function initPixelTracking(tracker: BaseTracker) { data.map((event) => event && typeof event === "object" ? sendOnePixel(eventType, event as Record) - : Promise.resolve({ attempts: 0, success: false }) + : Promise.resolve({ attempts: 1, success: false }) ) ); return { - attempts: Math.max(0, ...results.map((result) => result.attempts)), + attempts: Math.max(1, ...results.map((result) => result.attempts)), success: results.every((result) => result.success), }; } if (typeof data !== "object" || data === null) { - return { attempts: 0, success: false }; + return { attempts: 1, success: false }; } return sendOnePixel(eventType, data as Record); }; @@ -158,10 +249,8 @@ export function initPixelTracking(tracker: BaseTracker) { }; }; - tracker.sendBeacon = (data: unknown, endpoint = "/") => { - sendToPixel(endpoint, data); - return true; - }; - - tracker.sendBatchBeacon = () => false; + // Image loads cannot synchronously prove remote acceptance. Returning false + // keeps BaseTracker's queue for its fetch-style pixel fallback instead of + // dropping it immediately on an unverified load. + tracker.sendBeacon = () => false; } diff --git a/packages/tracker/tests/unit/pixel.test.ts b/packages/tracker/tests/unit/pixel.test.ts index e6c9e3e25..dc4a74ae6 100644 --- a/packages/tracker/tests/unit/pixel.test.ts +++ b/packages/tracker/tests/unit/pixel.test.ts @@ -1,64 +1,162 @@ -import { afterEach, describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, jest, mock, test } from "bun:test"; import type { BaseTracker } from "../../src/core/tracker"; +import type { TrackerOptions } from "../../src/core/types"; import { initPixelTracking } from "../../src/plugins/pixel"; -const originalImage = globalThis.Image; +const originalImage = Object.getOwnPropertyDescriptor(globalThis, "Image"); afterEach(() => { - globalThis.Image = originalImage; + if (jest.isFakeTimers()) { + jest.clearAllTimers(); + jest.useRealTimers(); + } + if (originalImage) { + Object.defineProperty(globalThis, "Image", originalImage); + return; + } + Reflect.deleteProperty(globalThis, "Image"); }); -function trackerFixture(): BaseTracker { - const tracker = { +function createTracker( + overrides: Partial = {} +): BaseTracker { + return { options: { apiUrl: "https://basket.example", clientId: "ws_test", - enableBatching: true, enableRetries: true, initialRetryDelay: 0, maxRetries: 1, sdk: "web", sdkVersion: "test", + ...overrides, + }, + api: { + cancelPendingRequests: mock(() => {}), }, - api: {}, - sendBatchBeacon: () => false, - sendBeacon: () => false, } as unknown as BaseTracker; - initPixelTracking(tracker); - return tracker; } -describe("pixel tracking", () => { - test("retries a failed image load", async () => { - const requests: string[] = []; - - globalThis.Image = class { - onerror: (() => void) | null = null; - onload: (() => void) | null = null; - - set src(value: string) { - requests.push(value); - queueMicrotask(() => { - if (requests.length === 1) { - this.onerror?.(); - return; - } - this.onload?.(); - }); +function installImageOutcomes( + outcomes: Array<"error" | "load" | "pending"> +): string[] { + const requests: string[] = []; + + class MockImage { + onerror: (() => void) | null = null; + onload: (() => void) | null = null; + + set src(value: string) { + requests.push(value); + if (!value) { + return; + } + const outcome = outcomes.shift() ?? "error"; + if (outcome === "pending") { + return; } - } as unknown as typeof Image; + queueMicrotask(() => { + if (outcome === "load") { + this.onload?.(); + return; + } + this.onerror?.(); + }); + } + } + + Object.defineProperty(globalThis, "Image", { + configurable: true, + value: MockImage, + }); + return requests; +} + + +describe("pixel transport", () => { + test("retries an unacknowledged pixel load with the same event identity", async () => { + const requests = installImageOutcomes(["error", "load"]); + const tracker = createTracker(); + initPixelTracking(tracker); + + const result = await tracker.api.fetch("/", { + eventId: "event_1", + name: "pageview", + }); + + expect(result).toMatchObject({ ok: true, attempts: 2 }); + expect(requests).toHaveLength(2); + const first = new URL(requests[0] ?? ""); + const second = new URL(requests[1] ?? ""); + expect(first.searchParams.get("eventId")).toBe("event_1"); + expect(second.searchParams.get("eventId")).toBe("event_1"); + }); + + test("does not report an unverified image load as a beacon success", () => { + const tracker = createTracker(); + initPixelTracking(tracker); - const tracker = trackerFixture(); - const result = await tracker.api.fetch("/track", { - eventId: "evt_1", - name: "signup_completed", + expect(tracker.sendBeacon({ eventId: "event_1" }, "/")).toBe(false); + }); + + test("cancels an active image request when tracking is cleared", async () => { + const requests = installImageOutcomes(["pending"]); + const tracker = createTracker(); + initPixelTracking(tracker); + + const delivery = tracker.api.fetch("/", { eventId: "event_1" }); + tracker.api.cancelPendingRequests(); + + expect(await delivery).toMatchObject({ + ok: false, + attempts: 1, }); + expect(requests).toHaveLength(2); + expect(requests[1]).toBe(""); + }); + + test("bounds a pixel image load that never completes", async () => { + jest.useFakeTimers(); + const requests = installImageOutcomes(["pending"]); + const tracker = createTracker({ maxRetries: 0 }); + initPixelTracking(tracker); - expect(result).toMatchObject({ - attempts: 2, - ok: true, + const delivery = tracker.api.fetch("/", { eventId: "event_1" }); + jest.advanceTimersByTime(10_000); + + expect(await delivery).toMatchObject({ + ok: false, + attempts: 1, }); expect(requests).toHaveLength(2); - expect(new URL(requests[0] ?? "").pathname).toBe("/px.jpg"); + expect(requests[1]).toBe(""); + }); + + test("cancels a scheduled retry without waiting for the backoff", async () => { + jest.useFakeTimers(); + const requests = installImageOutcomes(["error"]); + const tracker = createTracker({ initialRetryDelay: 30_000 }); + initPixelTracking(tracker); + + const delivery = tracker.api.fetch("/", { eventId: "event_1" }); + await Promise.resolve(); + tracker.api.cancelPendingRequests(); + + expect(await delivery).toMatchObject({ + ok: false, + attempts: 1, + }); + expect(requests).toHaveLength(1); + }); + + test("uses the retry fallback for an invalid maxRetries value", async () => { + const requests = installImageOutcomes(["error", "error", "error", "error"]); + const tracker = createTracker({ maxRetries: Number.NaN }); + initPixelTracking(tracker); + + const result = await tracker.api.fetch("/", { eventId: "event_1" }); + + expect(result).toMatchObject({ ok: false, attempts: 4 }); + expect(requests).toHaveLength(4); }); });