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.
8.9 kB ยท 302 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303import { getLogger } from "@logtape/logtape";import { and, count, db, eq, gte, inArray, schema } from "@openstatus/db";import type { Incident, MonitorStatus } from "@openstatus/db/src/schema";import { selectMonitorSchema, selectNotificationSchema, selectWorkspaceSchema,} from "@openstatus/db/src/schema";import type { Region } from "@openstatus/db/src/schema/constants";import { Effect, Schedule } from "effect";
import { checkerAudit } from "../utils/audit-log";import { providerToFunction } from "./utils";
const logger = getLogger("workflow");
export const triggerNotifications = async ({ monitorId, statusCode, message, notifType, cronTimestamp, incidentId, regions, latency,}: { monitorId: string; statusCode?: number; message?: string; notifType: "alert" | "recovery" | "degraded"; cronTimestamp: number; incidentId?: number; regions?: string[]; latency?: number;}): Promise<{ notificationId: number; provider: string }[]> => { logger.info("Triggering alerting", { monitor_id: monitorId, notification_type: notifType, });
const triggered: { notificationId: number; provider: string }[] = [];
let incident: Incident | undefined; if (incidentId) { try { incident = await db.query.incidentTable.findFirst({ where: eq(schema.incidentTable.id, incidentId), }); } catch (err) { logger.warn("Failed to fetch incident data", { incident_id: incidentId, error_message: err instanceof Error ? err.message : String(err), }); } }
const notifications = await db .select() .from(schema.notificationsToMonitors) .innerJoin( schema.notification, eq(schema.notification.id, schema.notificationsToMonitors.notificationId), ) .innerJoin( schema.monitor, eq(schema.monitor.id, schema.notificationsToMonitors.monitorId), ) .where(eq(schema.monitor.id, Number(monitorId))) .all(); for (const notif of notifications) { // for sms check we are in the quota if (notif.notification.provider === "sms") { if (notif.notification.workspaceId === null) { continue; }
const workspace = await db .select() .from(schema.workspace) .where(eq(schema.workspace.id, notif.notification.workspaceId));
if (workspace.length !== 1) { continue; }
const data = selectWorkspaceSchema.parse(workspace[0]);
const oneMonthAgo = new Date(); oneMonthAgo.setMonth(oneMonthAgo.getMonth() - 1);
const smsNotification = await db .select() .from(schema.notification) .where( and( eq(schema.notification.workspaceId, notif.notification.workspaceId), eq(schema.notification.provider, "sms"), ), ); const ids = smsNotification.map((notification) => notification.id);
const smsSent = await db .select({ count: count() }) .from(schema.notificationTrigger) .where( and( gte( schema.notificationTrigger.cronTimestamp, Math.floor(oneMonthAgo.getTime() / 1000), ), inArray(schema.notificationTrigger.notificationId, ids), ), ) .all();
if ((smsSent[0]?.count ?? 0) > data.limits["sms-limit"]) { logger.warn( `SMS quota exceeded for workspace ${notif.notification.workspaceId}`, ); continue; } } logger.info("Sending notification", { monitor_id: monitorId, provider: notif.notification.provider, notification_type: notifType, notification_id: notif.notification.id, }); const monitor = selectMonitorSchema.parse(notif.monitor); try { await insertNotificationTrigger({ monitorId: monitor.id, notificationId: notif.notification.id, cronTimestamp: cronTimestamp, }); } catch (_e) { logger.error("notification trigger already exists dont send again"); continue; } triggered.push({ notificationId: notif.notification.id, provider: notif.notification.provider, }); switch (notifType) { case "alert": const alertResult = Effect.tryPromise({ try: () => providerToFunction[notif.notification.provider].sendAlert({ monitor, notification: selectNotificationSchema.parse(notif.notification), statusCode, message, incident, cronTimestamp, regions, latency, }),
catch: (_unknown) => new Error( `Failed sending notification via ${notif.notification.provider} for monitor ${monitorId}`, ), }).pipe( Effect.retry({ times: 3, schedule: Schedule.exponential("1000 millis"), }), ); await Effect.runPromise(alertResult).catch((err) => logger.error("Failed to send alert notification", { monitor_id: monitorId, provider: notif.notification.provider, error_message: err instanceof Error ? err.message : String(err), }), ); break; case "recovery": const recoveryResult = Effect.tryPromise({ try: () => providerToFunction[notif.notification.provider].sendRecovery({ monitor, notification: selectNotificationSchema.parse(notif.notification), statusCode, message, incident, cronTimestamp, regions, latency, }), catch: (_unknown) => new Error( `Failed sending notification via ${notif.notification.provider} for monitor ${monitorId}`, ), }).pipe( Effect.retry({ times: 3, schedule: Schedule.exponential("1000 millis"), }), ); await Effect.runPromise(recoveryResult).catch((err) => logger.error("Failed to send recovery notification", { monitor_id: monitorId, provider: notif.notification.provider, error_message: err instanceof Error ? err.message : String(err), }), ); break; case "degraded": const degradedResult = Effect.tryPromise({ try: () => providerToFunction[notif.notification.provider].sendDegraded({ monitor, notification: selectNotificationSchema.parse(notif.notification), statusCode, message, incident, cronTimestamp, regions, latency, }), catch: (_unknown) => new Error( `Failed sending notification via ${notif.notification.provider} for monitor ${monitorId}`, ), }).pipe( Effect.retry({ times: 3, schedule: Schedule.exponential("1000 millis"), }), ); await Effect.runPromise(degradedResult).catch((err) => logger.error("Failed to send degraded notification", { monitor_id: monitorId, provider: notif.notification.provider, error_message: err instanceof Error ? err.message : String(err), }), ); break; } // ALPHA await checkerAudit.publishAuditLog({ id: `monitor:${monitorId}`, action: "notification.sent", targets: [{ id: monitorId, type: "monitor" }], metadata: { provider: notif.notification.provider, cronTimestamp, type: notifType, notificationId: notif.notification.id, }, }); }
return triggered;};
const insertNotificationTrigger = async ({ monitorId, notificationId, cronTimestamp,}: { monitorId: number; notificationId: number; cronTimestamp: number;}) => { await db .insert(schema.notificationTrigger) .values({ monitorId: Number(monitorId), notificationId: notificationId, cronTimestamp: cronTimestamp, }) .returning();};
export const upsertMonitorStatus = async ({ monitorId, status, region,}: { monitorId: string; status: MonitorStatus; region: Region;}) => { const newData = await db .insert(schema.monitorStatusTable) .values({ status, region, monitorId: Number(monitorId) }) .onConflictDoUpdate({ target: [ schema.monitorStatusTable.monitorId, schema.monitorStatusTable.region, ], set: { status, updatedAt: new Date() }, }) .returning(); logger.debug("Upserted monitor status", { monitor_id: monitorId, region, status, updated_at: newData[0]?.updatedAt, });};