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.
16 kB ยท 494 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495import { getLogger } from "@logtape/logtape";import { and, db, eq, gte, inArray, isNull, lte, schema, sql,} from "@openstatus/db";import { monitorStatusSchema } from "@openstatus/db/src/schema/monitors/validation";import type { Context } from "hono";import { z } from "zod";
import { env } from "../env";import type { Env } from "../index";import { checkerAudit } from "../utils/audit-log";import { triggerNotifications } from "./alerting";import { findOpenIncident, resolveIncident } from "./incident-utils";
const logger = getLogger(["workflow"]);
const payloadSchema = z.object({ monitorId: z.string(), privateLocationId: z.string(), status: monitorStatusSchema, cronTimestamp: z.number(), message: z.string().optional(), statusCode: z.number().optional(), latency: z.number().optional(),});
export async function updateStatusPrivate(c: Context<Env>) { const auth = c.req.header("Authorization"); if (auth !== `Basic ${env().CRON_SECRET}`) { logger.error("Unauthorized"); return c.text("Unauthorized", 401); }
const result = payloadSchema.safeParse(await c.req.json()); if (!result.success) { return c.text("Unprocessable Entity", 422); }
const { monitorId, privateLocationId, status, cronTimestamp, message, statusCode, latency, } = result.data;
const event = c.get("event"); const monitorIdNumber = Number(monitorId); const privateLocationIdNumber = Number(privateLocationId);
// Set event data for OTel/Axiom structured logging (matches cloud checker pattern) // Note: Private location pings are ingested to Tinybird separately via Go code in // apps/private-location/internal/tinybird/client.go, not through this event object if (event) { event.status_update = { status: status, message: message, region: privateLocationId, // Private location ID as region string status_code: statusCode, cron_timestamp: cronTimestamp, latency_ms: latency, monitorId: monitorIdNumber, }; }
try { const monitor = await db .select() .from(schema.monitor) .where(eq(schema.monitor.id, monitorIdNumber)) .get();
if (!monitor || monitor.deletedAt || !monitor.active) { return c.json({ success: true }, 200); }
const now = new Date(); const activeMaintenance = await db .select({ id: schema.maintenance.id }) .from(schema.maintenance) .innerJoin( schema.maintenancesToPageComponents, eq( schema.maintenancesToPageComponents.maintenanceId, schema.maintenance.id, ), ) .innerJoin( schema.pageComponent, eq( schema.pageComponent.id, schema.maintenancesToPageComponents.pageComponentId, ), ) .where( and( lte(schema.maintenance.from, now), gte(schema.maintenance.to, now), eq(schema.pageComponent.monitorId, monitorIdNumber), ), ) .get();
if (activeMaintenance) { return c.json({ success: true }, 200); }
const attachment = await db .select({ name: schema.privateLocation.name }) .from(schema.privateLocationToMonitors) .innerJoin( schema.privateLocation, eq( schema.privateLocation.id, schema.privateLocationToMonitors.privateLocationId, ), ) .where( and( eq(schema.privateLocationToMonitors.monitorId, monitorIdNumber), eq( schema.privateLocationToMonitors.privateLocationId, privateLocationIdNumber, ), isNull(schema.privateLocationToMonitors.deletedAt), ), ) .get();
if (!attachment) { return c.json({ success: true }, 200); }
const priorRow = await db .select({ status: schema.privateLocationMonitorStatus.status }) .from(schema.privateLocationMonitorStatus) .where( and( eq(schema.privateLocationMonitorStatus.monitorId, monitorIdNumber), eq( schema.privateLocationMonitorStatus.privateLocationId, privateLocationIdNumber, ), ), ) .get();
const priorStatus = priorRow?.status ?? "active";
const upserted = await db .insert(schema.privateLocationMonitorStatus) .values({ monitorId: monitorIdNumber, privateLocationId: privateLocationIdNumber, status, cronTimestamp, }) .onConflictDoUpdate({ target: [ schema.privateLocationMonitorStatus.monitorId, schema.privateLocationMonitorStatus.privateLocationId, ], set: { status, cronTimestamp, updatedAt: new Date() }, setWhere: sql`excluded.cron_timestamp > ${schema.privateLocationMonitorStatus.cronTimestamp}`, }) .returning();
if (upserted.length === 0 || status === priorStatus) { return c.json({ success: true }, 200); }
const regions = [attachment.name];
// Check if monitor has cloud regions const hasCloudRegions = monitor.regions && monitor.regions.trim().length > 0;
// Query all private locations for threshold check const allLocations = await db .select({ id: schema.privateLocationToMonitors.privateLocationId }) .from(schema.privateLocationToMonitors) .where( and( eq(schema.privateLocationToMonitors.monitorId, monitorIdNumber), isNull(schema.privateLocationToMonitors.deletedAt), ), ) .all();
const numberOfLocations = allLocations.length; const locationIds = allLocations .map((loc) => loc.id) .filter((id): id is number => id !== null);
// Count how many locations report this status const locationsWithStatus = await db .select({ privateLocationId: schema.privateLocationMonitorStatus.privateLocationId, }) .from(schema.privateLocationMonitorStatus) .where( and( eq(schema.privateLocationMonitorStatus.monitorId, monitorIdNumber), eq(schema.privateLocationMonitorStatus.status, status), inArray( schema.privateLocationMonitorStatus.privateLocationId, locationIds, ), ), ) .all();
const affectedLocationCount = locationsWithStatus.length;
// Apply โฅ50% threshold (matching cloud checker logic) const shouldTriggerIncident = !hasCloudRegions && (affectedLocationCount >= numberOfLocations / 2 || numberOfLocations === 1);
let incident = null; let triggeredNotifications: { notificationId: number; provider: string }[] = [];
switch (status) { case "error": // Update monitor status for private-only monitors if ( !hasCloudRegions && shouldTriggerIncident && monitor.status !== "error" ) { logger.info("Monitor status changed to error", { monitor_id: monitor.id, workspace_id: monitor.workspaceId, }); await db .update(schema.monitor) .set({ status: "error" }) .where(eq(schema.monitor.id, monitorIdNumber)); }
// Create incident only if private-only monitor AND threshold met if (shouldTriggerIncident) { try { const existingIncident = await findOpenIncident(monitorIdNumber); if (!existingIncident) { const [newIncident] = await db .insert(schema.monitorIncidentTable) .values({ monitorId: monitorIdNumber, workspaceId: monitor.workspaceId, startedAt: new Date(cronTimestamp), }) .returning();
if (newIncident?.id) { incident = newIncident; await checkerAudit.publishAuditLog({ id: `monitor:${monitorId}`, action: "incident.created", targets: [{ id: monitorId, type: "monitor" }], metadata: { cronTimestamp, incidentId: newIncident.id }, }); logger.info("Created incident", { incident_id: newIncident.id, monitor_id: monitorId, affected_location_count: affectedLocationCount, total_locations: numberOfLocations, }); } } else { incident = existingIncident; logger.info("Already in incident", { incident_id: existingIncident.id, }); } } catch (error) { // Check if this is a constraint violation (race condition) const errorMessage = error instanceof Error ? error.message : String(error); if ( errorMessage.includes("UNIQUE constraint") || errorMessage.includes("unique") ) { // Another request created the incident concurrently, fetch it logger.info( "Concurrent incident creation detected, fetching existing", { monitor_id: monitorId, }, ); const existingIncident = await findOpenIncident(monitorIdNumber); if (existingIncident) { incident = existingIncident; } } else { logger.error("Failed to create incident", { monitor_id: monitorId, error_message: errorMessage, }); } } }
await checkerAudit.publishAuditLog({ id: `monitor:${monitorId}`, action: "monitor.failed", targets: [{ id: monitorId, type: "monitor" }], metadata: { region: privateLocationId, statusCode: statusCode ?? -1, message, cronTimestamp, latency, }, }); // Trigger notifications for cloud monitors always, private-only when threshold met AND status changed if ( hasCloudRegions || (shouldTriggerIncident && monitor.status !== "error") ) { triggeredNotifications = await triggerNotifications({ monitorId, statusCode, message, notifType: "alert", cronTimestamp, regions, latency, incidentId: incident?.id, }); } break; case "degraded": // Update monitor status for private-only monitors if ( !hasCloudRegions && shouldTriggerIncident && monitor.status !== "degraded" ) { logger.info("Monitor status changed to degraded", { monitor_id: monitor.id, workspace_id: monitor.workspaceId, }); await db .update(schema.monitor) .set({ status: "degraded" }) .where(eq(schema.monitor.id, monitorIdNumber)); }
// Resolve incident if private-only monitor AND threshold met if (shouldTriggerIncident) { try { const incidents = await resolveIncident({ monitorId, cronTimestamp, }); incident = incidents[0] ?? null; } catch (error) { logger.warning( "Failed to resolve incident on degraded transition", { monitor_id: monitorId, error_message: error instanceof Error ? error.message : String(error), }, ); // Continue with notifications even if resolution fails } }
await checkerAudit.publishAuditLog({ id: `monitor:${monitorId}`, action: "monitor.degraded", targets: [{ id: monitorId, type: "monitor" }], metadata: { region: privateLocationId, statusCode: statusCode ?? -1, cronTimestamp, latency, }, }); // Trigger notifications for cloud monitors always, private-only when threshold met AND status changed if ( hasCloudRegions || (shouldTriggerIncident && monitor.status !== "degraded") ) { triggeredNotifications = await triggerNotifications({ monitorId, statusCode, message, notifType: "degraded", cronTimestamp, regions, latency, incidentId: incident?.id, }); } break; case "active": // Update monitor status for private-only monitors if ( !hasCloudRegions && shouldTriggerIncident && monitor.status !== "active" ) { logger.info("Monitor status changed to active", { monitor_id: monitor.id, workspace_id: monitor.workspaceId, }); await db .update(schema.monitor) .set({ status: "active" }) .where(eq(schema.monitor.id, monitorIdNumber)); }
// Resolve incident if private-only monitor AND threshold met if (shouldTriggerIncident) { try { const incidents = await resolveIncident({ monitorId, cronTimestamp, }); incident = incidents[0] ?? null; } catch (error) { logger.warning("Failed to resolve incident on active transition", { monitor_id: monitorId, error_message: error instanceof Error ? error.message : String(error), }); // Continue with notifications even if resolution fails } }
await checkerAudit.publishAuditLog({ id: `monitor:${monitorId}`, action: "monitor.recovered", targets: [{ id: monitorId, type: "monitor" }], metadata: { region: privateLocationId, statusCode: statusCode ?? -1, cronTimestamp, latency, }, }); // Trigger notifications for cloud monitors always, private-only when threshold met AND status changed if ( hasCloudRegions || (shouldTriggerIncident && monitor.status !== "active") ) { triggeredNotifications = await triggerNotifications({ monitorId, statusCode, message, notifType: "recovery", cronTimestamp, regions, latency, incidentId: incident?.id, }); } break; }
// Add notification outcome to event for OTel logging if (event?.status_update) { (event.status_update as Record<string, unknown>).notificationTriggered = triggeredNotifications.length > 0; (event.status_update as Record<string, unknown>).notifications = triggeredNotifications; }
return c.json({ success: true }, 200); } catch (error) { logger.error("Failed to update private location status", { monitor_id: monitorId, private_location_id: privateLocationId, error_message: error instanceof Error ? error.message : String(error), }); return c.text("Internal Server Error", 500); }}