Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/openstatusHQ/openstatus. ๐ซ Status page with uptime monitoring & API monitoring as code ๐ซ openstatus.dev
bun drizzle-orm monitoring monitoring-as-code nextjs observability on-call open-source shadcn-ui status-page statuspage synthetic-monitoring tinybird turso uptime uptime-checker uptime-monitor
Something went wrong. Try again.
22 kB ยท 794 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795import { Events } from "@openstatus/analytics";import { deserialize, dnsRecords, headerAssertion, jsonBodyAssertion, recordAssertion, statusAssertion, textBodyAssertion,} from "@openstatus/assertions";import { and, db, eq } from "@openstatus/db";import { monitor, selectMonitorSchema } from "@openstatus/db/src/schema";import { monitorRegionSchema } from "@openstatus/db/src/schema/constants";import { type httpPayloadSchema, type grpcPayloadSchema, type icmpPayloadSchema, GRPC_TLS_MODES, headerPairSchema, safeUrlSchema, type tpcPayloadSchema, transformHeaders,} from "@openstatus/utils";import { TRPCError } from "@trpc/server";import { z } from "zod";
import { env } from "../env";import { createTRPCRouter, protectedProcedure } from "../trpc";
const ABORT_TIMEOUT = 10000;const CHECKER_BASE_URL = ( env.CHECKER_URL || "https://openstatus-checker.fly.dev").replace(/\/+$/, "");
// PingICMP treats its timeout as the deadline for the whole check, so omitting// it means a deadline of "now": the send loop breaks before the first packet// and every test reports "no reply". Kept under ABORT_TIMEOUT so the checker// answers before the fetch above gives up.const ICMP_TEST_TIMEOUT = 5000;
// Kept under ABORT_TIMEOUT so the checker answers before the fetch gives up.const GRPC_TEST_TIMEOUT = 5000;
// Unreachable targets, failed assertions and timeouts are expected outcomes// already surfaced to the user as BAD_REQUEST; only unexpected failures should// reach the logs (and Sentry).function toCheckerError( error: unknown, label: string, fallback: string,): TRPCError { if (error instanceof TRPCError) { if (error.code !== "BAD_REQUEST") { console.error(`Checker ${label} test failed`, error); } return error; }
if (error instanceof Error && error.name === "TimeoutError") { return new TRPCError({ code: "BAD_REQUEST", message: `The ${label} check did not complete within ${ ABORT_TIMEOUT / 1000 } seconds. Please try again.`, }); }
console.error(`Checker ${label} test failed`, error); return new TRPCError({ code: "INTERNAL_SERVER_ERROR", message: fallback });}
// Input schemasconst httpTestInput = z.object({ url: safeUrlSchema, method: z .enum([ "GET", "HEAD", "OPTIONS", "POST", "PUT", "DELETE", "PATCH", "CONNECT", "TRACE", ]) .prefault("GET"), headers: z.array(headerPairSchema).optional(), body: z.string().optional(), region: monitorRegionSchema.prefault("ams"), assertions: z .array( z.discriminatedUnion("type", [ statusAssertion, headerAssertion, textBodyAssertion, jsonBodyAssertion, recordAssertion, ]), ) .prefault([]),});
const tcpTestInput = z.object({ url: z.string(), region: monitorRegionSchema.prefault("ams"),});
const dnsTestInput = z.object({ url: z.string(), region: monitorRegionSchema.prefault("ams"), assertions: z .array( z.discriminatedUnion("type", [ recordAssertion, statusAssertion, headerAssertion, textBodyAssertion, jsonBodyAssertion, ]), ) .prefault([]),});
const icmpTestInput = z.object({ url: z.string(), region: monitorRegionSchema.prefault("ams"),});
const grpcTestInput = z.object({ url: z.string(), service: z.string().optional(), tls: z.enum(GRPC_TLS_MODES).prefault("tls"), headers: z.array(headerPairSchema).optional(), region: monitorRegionSchema.prefault("ams"),});
export const grpcOutput = z .object({ state: z.literal("success").prefault("success"), type: z.literal("grpc").prefault("grpc"), jobType: z.literal("grpc").optional(), requestId: z.number().optional(), workspaceId: z.number().optional(), monitorId: z.number().optional(), timestamp: z.number(), timing: z.object({ dnsStart: z.number(), dnsDone: z.number(), connectStart: z.number(), connectDone: z.number(), tlsHandshakeStart: z.number(), tlsHandshakeDone: z.number(), firstByteStart: z.number(), firstByteDone: z.number(), transferStart: z.number(), transferDone: z.number(), }), latency: z.number().optional(), servingStatus: z.string().optional(), service: z.string().optional(), grpcCode: z.number().optional(), completed: z.boolean().optional(), errorMessage: z.string().optional(), error: z.number().optional(), region: monitorRegionSchema, }) .or( z.object({ state: z.literal("error").prefault("error"), message: z.string(), }), );
export const icmpOutput = z .object({ state: z.literal("success").prefault("success"), type: z.literal("icmp").prefault("icmp"), requestId: z.number().optional(), workspaceId: z.number().optional(), monitorId: z.number().optional(), timestamp: z.number(), timing: z.object({ rtts: z.array(z.number()), }), latency: z.number().optional(), latencyMin: z.number().optional(), latencyMax: z.number().optional(), packetsSent: z.number().optional(), packetsReceived: z.number().optional(), error: z.string().optional(), region: monitorRegionSchema, }) .or( z.object({ state: z.literal("error").prefault("error"), message: z.string(), }), );
export const tcpOutput = z .object({ state: z.literal("success").prefault("success"), type: z.literal("tcp").prefault("tcp"), requestId: z.number().optional(), workspaceId: z.number().optional(), monitorId: z.number().optional(), timestamp: z.number(), timing: z.object({ tcpStart: z.number(), tcpDone: z.number(), }), error: z.string().optional(), region: monitorRegionSchema, latency: z.number().optional(), }) .or( z.object({ state: z.literal("error").prefault("error"), message: z.string(), }), );
export const httpOutput = z .object({ state: z.literal("success").prefault("success"), type: z.literal("http").prefault("http"), status: z.number(), latency: z.number(), headers: z.record(z.string(), z.string()), timestamp: z.number(), timing: z.object({ dnsStart: z.number(), dnsDone: z.number(), connectStart: z.number(), connectDone: z.number(), tlsHandshakeStart: z.number(), tlsHandshakeDone: z.number(), firstByteStart: z.number(), firstByteDone: z.number(), transferStart: z.number(), transferDone: z.number(), }), body: z.string().optional().nullable(), region: monitorRegionSchema, }) .or( z.object({ state: z.literal("error").prefault("error"), message: z.string(), }), ) .or( // A target the checker could not reach (timeout, DNS, refused): `checker.Http` // answers 200 with `error` set and `status`/`headers` omitted. z .object({ error: z.string().min(1), timestamp: z.number() }) .transform(({ error }) => ({ state: "error" as const, message: error })), );
export const dnsOutput = z .object({ state: z.literal("success").prefault("success"), type: z.literal("dns").prefault("dns"), records: z .partialRecord(z.enum(dnsRecords), z.array(z.string())) .prefault({}), latency: z.number().optional(), timestamp: z.number(), region: monitorRegionSchema, }) .or( z.object({ state: z.literal("error").prefault("error"), message: z.string(), }), );
export async function testHttp(input: z.infer<typeof httpTestInput>) { // Reject requests to our own domain to avoid loops if (input.url.includes("openstatus.dev")) { throw new TRPCError({ code: "BAD_REQUEST", message: "Self-requests are not allowed", }); }
try { const res = await fetch(`${CHECKER_BASE_URL}/ping/${input.region}`, { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": input.region, }, body: JSON.stringify({ url: input.url, method: input.method, headers: input.headers?.reduce( (acc, { key, value }) => { if (!key) return acc; return { ...acc, [key]: value }; }, {} as Record<string, string>, ), body: input.body, }), signal: AbortSignal.timeout(ABORT_TIMEOUT), });
const json = await res.json(); const result = httpOutput.safeParse(json);
if (!result.success) { console.error( `Checker HTTP test failed for ${input.url}:`, result.error.message, ); throw new TRPCError({ code: "BAD_REQUEST", message: "Checker response is not valid. Please try again. If the problem persists, please contact support.", }); }
if (result.data.state === "error") { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.message, }); }
if (result.data.state === "success") { const { body, headers, status } = result.data;
const assertions = deserialize(JSON.stringify(input.assertions)).map( (assertion) => assertion.assert({ body: body ?? "", header: headers ?? {}, status: status, }), );
if (assertions.some((assertion) => !assertion.success)) { throw new TRPCError({ code: "BAD_REQUEST", message: `Assertion error: ${ assertions.find((assertion) => !assertion.success)?.message }`, }); }
if (assertions.length === 0 && (status < 200 || status >= 300)) { throw new TRPCError({ code: "BAD_REQUEST", message: `Assertion error: The response status was not 2XX: ${status}.`, }); } }
return result.data; } catch (error) { throw toCheckerError( error, "HTTP", error instanceof Error ? error.message : "HTTP check failed", ); }}
export async function testTcp(input: z.infer<typeof tcpTestInput>) { try { const res = await fetch(`${CHECKER_BASE_URL}/tcp/${input.region}`, { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": input.region, }, body: JSON.stringify({ uri: input.url }), signal: AbortSignal.timeout(ABORT_TIMEOUT), });
const json = await res.json(); const result = tcpOutput.safeParse(json);
if (!result.success) { console.error( `Checker TCP test failed for ${input.url}:`, result.error.message, ); throw new TRPCError({ code: "BAD_REQUEST", message: `Checker response is not valid. Please try again. If the problem persists, please contact support. ${result.error.message}`, }); }
if (result.data.state === "error") { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.message, }); }
return result.data; } catch (error) { throw toCheckerError(error, "TCP", "TCP check failed"); }}
export async function testDns(input: z.infer<typeof dnsTestInput>) { try { const res = await fetch(`${CHECKER_BASE_URL}/dns/${input.region}`, { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": input.region, }, body: JSON.stringify({ uri: input.url, }), signal: AbortSignal.timeout(ABORT_TIMEOUT), });
const json = await res.json(); const result = dnsOutput.safeParse(json);
if (!result.success) { console.error( `Checker DNS test failed for ${input.url}:`, result.error.message, ); throw new TRPCError({ code: "BAD_REQUEST", message: `Checker response is not valid. Please try again. If the problem persists, please contact support. ${result.error.message}`, }); }
if (result.data.state === "error") { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.message, }); }
if (result.data.state === "success") { const { records } = result.data;
const assertions = deserialize(JSON.stringify(input.assertions)).map( (assertion) => assertion.assert({ records }), );
if (assertions.some((assertion) => !assertion.success)) { throw new TRPCError({ code: "BAD_REQUEST", message: `Assertion error: ${ assertions.find((assertion) => !assertion.success)?.message }`, }); } }
return result.data; } catch (error) { throw toCheckerError(error, "DNS", "DNS check failed"); }}
export async function testIcmp(input: z.infer<typeof icmpTestInput>) { try { const res = await fetch(`${CHECKER_BASE_URL}/icmp/${input.region}`, { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": input.region, }, body: JSON.stringify({ uri: input.url, timeout: ICMP_TEST_TIMEOUT, }), signal: AbortSignal.timeout(ABORT_TIMEOUT), });
const json = await res.json(); const result = icmpOutput.safeParse(json);
if (!result.success) { console.error( `Checker ICMP test failed for ${input.url}:`, result.error.message, ); throw new TRPCError({ code: "BAD_REQUEST", message: `Checker response is not valid. Please try again. If the problem persists, please contact support. ${result.error.message}`, }); }
if (result.data.state === "error") { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.message, }); }
return result.data; } catch (error) { throw toCheckerError(error, "ICMP", "ICMP check failed"); }}
export async function testGrpc(input: z.infer<typeof grpcTestInput>) { try { const res = await fetch(`${CHECKER_BASE_URL}/grpc/${input.region}`, { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": input.region, }, body: JSON.stringify({ uri: input.url, service: input.service, tls: input.tls, headers: transformHeaders(input.headers ?? []), timeout: GRPC_TEST_TIMEOUT, }), signal: AbortSignal.timeout(ABORT_TIMEOUT), });
const json = await res.json(); const result = grpcOutput.safeParse(json);
if (!result.success) { console.error( `Checker gRPC test failed for ${input.url}:`, result.error.message, ); throw new TRPCError({ code: "BAD_REQUEST", message: `Checker response is not valid. Please try again. If the problem persists, please contact support. ${result.error.message}`, }); }
if (result.data.state === "error") { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.message, }); }
// Only a transport failure comes back as `state: "error"`. An RPC that // completed but answered NOT_SERVING / SERVICE_UNKNOWN โ or a server with no // health service at all โ returns the full response, where `state` is absent // and prefaults to "success". `error` is omitempty, so it is present only // when the check failed. Mirrors testHttp rejecting a non-2XX status: the // target is reachable, but saving it would create a monitor that is already // down. if (result.data.error === 1) { throw new TRPCError({ code: "BAD_REQUEST", message: result.data.errorMessage || `The health check did not report SERVING${ result.data.servingStatus ? `: ${result.data.servingStatus}` : "" }`, }); }
return result.data; } catch (error) { throw toCheckerError(error, "gRPC", "gRPC check failed"); }}
export async function triggerChecker( input: z.infer<typeof selectMonitorSchema>,) { let payload: | z.infer<typeof httpPayloadSchema> | z.infer<typeof tpcPayloadSchema> | z.infer<typeof icmpPayloadSchema> | z.infer<typeof grpcPayloadSchema> | null = null;
if (process.env.NODE_ENV !== "production") { return; }
const timestamp = Date.now();
if (input.jobType === "http") { payload = { workspaceId: String(input.workspaceId), monitorId: String(input.id), url: input.url, method: input.method || "GET", cronTimestamp: timestamp, body: input.body, headers: input.headers, status: "active", assertions: input.assertions ? JSON.parse(input.assertions) : null, degradedAfter: input.degradedAfter, timeout: input.timeout, trigger: "cron", otelConfig: input.otelEndpoint ? { endpoint: input.otelEndpoint, headers: transformHeaders(input.otelHeaders), } : undefined, retry: input.retry || 3, followRedirects: input.followRedirects || true, }; } if (input.jobType === "tcp") { payload = { workspaceId: String(input.workspaceId), monitorId: String(input.id), uri: input.url, status: "active", assertions: input.assertions ? JSON.parse(input.assertions) : null, cronTimestamp: timestamp, degradedAfter: input.degradedAfter, timeout: input.timeout, trigger: "cron", retry: input.retry || 3, otelConfig: input.otelEndpoint ? { endpoint: input.otelEndpoint, headers: transformHeaders(input.otelHeaders), } : undefined, followRedirects: input.followRedirects || true, }; } if (input.jobType === "dns") { payload = { workspaceId: String(input.workspaceId), monitorId: String(input.id), uri: input.url, status: "active", assertions: input.assertions ? JSON.parse(input.assertions) : null, cronTimestamp: timestamp, degradedAfter: input.degradedAfter, timeout: input.timeout, trigger: "cron", retry: input.retry || 3, otelConfig: input.otelEndpoint ? { endpoint: input.otelEndpoint, headers: transformHeaders(input.otelHeaders), } : undefined, followRedirects: input.followRedirects || true, }; } if (input.jobType === "icmp") { payload = { workspaceId: String(input.workspaceId), monitorId: String(input.id), uri: input.url, status: "active", cronTimestamp: timestamp, degradedAfter: input.degradedAfter, timeout: input.timeout, trigger: "cron", retry: input.retry || 3, otelConfig: input.otelEndpoint ? { endpoint: input.otelEndpoint, headers: transformHeaders(input.otelHeaders), } : undefined, }; } if (input.jobType === "grpc") { payload = { workspaceId: String(input.workspaceId), monitorId: String(input.id), uri: input.url, service: input.grpcService ?? undefined, tls: input.grpcTls ?? "tls", headers: transformHeaders(input.headers), status: "active", cronTimestamp: timestamp, degradedAfter: input.degradedAfter, timeout: input.timeout, trigger: "cron", retry: input.retry || 3, otelConfig: input.otelEndpoint ? { endpoint: input.otelEndpoint, headers: transformHeaders(input.otelHeaders), } : undefined, }; } const allResult = [];
for (const region of input.regions) { const res = fetch(generateUrl({ row: input }), { method: "POST", headers: { Authorization: `Basic ${env.CRON_SECRET}`, "Content-Type": "application/json", "fly-prefer-region": region, }, body: JSON.stringify(payload), signal: AbortSignal.timeout(ABORT_TIMEOUT), }); allResult.push(res); }
await Promise.allSettled(allResult);}
function generateUrl({ row }: { row: z.infer<typeof selectMonitorSchema> }) { switch (row.jobType) { case "http": return `${CHECKER_BASE_URL}/checker/http?monitor_id=${row.id}`; case "tcp": return `${CHECKER_BASE_URL}/checker/tcp?monitor_id=${row.id}`; case "dns": return `${CHECKER_BASE_URL}/checker/dns?monitor_id=${row.id}`; case "icmp": return `${CHECKER_BASE_URL}/checker/icmp?monitor_id=${row.id}`; case "grpc": return `${CHECKER_BASE_URL}/checker/grpc?monitor_id=${row.id}`; default: throw new Error("Invalid jobType"); }}
export const checkerRouter = createTRPCRouter({ testHttp: protectedProcedure .meta({ track: Events.TestMonitor }) .input(httpTestInput) .mutation(async ({ input }) => { return testHttp(input); }),
testTcp: protectedProcedure .meta({ track: Events.TestMonitor }) .input(tcpTestInput) .mutation(async ({ input }) => { return testTcp(input); }), testDns: protectedProcedure .meta({ track: Events.TestMonitor }) .input(dnsTestInput) .mutation(async ({ input }) => { return testDns(input); }), testIcmp: protectedProcedure .meta({ track: Events.TestMonitor }) .input(icmpTestInput) .mutation(async ({ input }) => { return testIcmp(input); }),
testGrpc: protectedProcedure .meta({ track: Events.TestMonitor }) .input(grpcTestInput) .mutation(async ({ input }) => { return testGrpc(input); }),
triggerChecker: protectedProcedure .input(z.object({ id: z.number() })) .mutation(async (opts) => { const m = await db .select() .from(monitor) .where( and( eq(monitor.id, opts.input.id), eq(monitor.workspaceId, opts.ctx.workspace.id), ), ) .get(); if (!m) { throw new TRPCError({ code: "NOT_FOUND", message: "Monitor not found", }); } const input = selectMonitorSchema.parse(m);
return await triggerChecker(input); }),});