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 ยท 755 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756import { getLogger } from "@logtape/logtape";import { db } from "@openstatus/db";import { listExternalServices } from "@openstatus/services/external-service";import type { ExternalServiceRow } from "@openstatus/services/external-service";import { applyDetectedProvider } from "@openstatus/services/external-service";import { type UpsertExternalComponentInput, upsertExternalComponentsForService,} from "@openstatus/services/external-service-component";import { type UpsertExternalIncidentInput, upsertExternalIncidentsForService,} from "@openstatus/services/external-service-incident";import { FetchError, detectProvider, fetchers,} from "@openstatus/status-fetcher";import type { NormalizedComponent, NormalizedIncident, StatusFetcher, StatusPageEntry, StatusResult,} from "@openstatus/status-fetcher";import { OSTinybird } from "@openstatus/tinybird";import { Effect } from "effect";import type { Context } from "hono";
import { env } from "../env";import { reportBackgroundError, reportDetectionStory, reportDetectionWriteFailure, reportFetchFailure, runSentryCron,} from "../lib/sentry";import { clearProbeStamp, decideDetectionAction, isBlocked, isSuspicious, shouldProbe,} from "./external-status-detect";
const logger = getLogger(["workflow", "external-status"]);
const tb = new OSTinybird({ token: env().TINY_BIRD_API_KEY, baseUrl: env().TINYBIRD_URL, noop: env().TINYBIRD_NOOP,});
// 10 per phase ร 3 phases = peak 30 concurrent HTTP requests upstream; keeps// Atlassian/Incident.io CDNs comfortable while still parallelising heavily.const PHASE_CONCURRENCY = 10;
function toStatusPageEntry(row: ExternalServiceRow): StatusPageEntry { return { id: row.slug, name: row.name, url: row.url, status_page_url: row.statusPageUrl, provider: row.provider, industry: row.industry, description: row.description ?? undefined, api_config: row.apiConfig ?? undefined, };}
type Snapshot = { id: string; indicator: string; status: string; status_message: string; fetched_at: number; updated_at: number; time_zone: string;};
function buildSnapshot(args: { entry: StatusPageEntry; result: StatusResult; fetchedAt: number;}): Snapshot { const { entry, result, fetchedAt } = args; return { id: entry.id, indicator: result.severity, status: result.status, status_message: result.description, fetched_at: fetchedAt, updated_at: result.updated_at, time_zone: result.timezone ?? "", };}
function toUpsertInput( incident: NormalizedIncident,): UpsertExternalIncidentInput { return { providerIncidentId: incident.providerIncidentId, name: incident.name, status: incident.status, impact: incident.impact, shortlink: incident.shortlink, startedAt: incident.startedAt, createdAt: incident.createdAt, resolvedAt: incident.resolvedAt, affectedComponentIds: incident.affectedComponentIds, raw: incident.raw, };}
type ComponentSnapshot = { component_id: string; external_service_id: number; indicator: string; status: string; fetched_at: number;};
function toComponentUpsertInput( component: NormalizedComponent,): UpsertExternalComponentInput { return { upstreamComponentId: component.upstreamComponentId, name: component.name, description: component.description, groupName: component.groupName, position: component.position, indicator: component.severity, status: component.status, };}
type PhaseCounts = { successCount: number; failureCount: number; skippedCount: number; total: number;};
type StatusPhaseOutcome = | { kind: "ok"; snapshot: Snapshot } | { kind: "no-fetcher"; slug: string } | { kind: "fail"; slug: string; reason: string; error: FetchError };
type IncidentPhaseOutcome = | { kind: "ok"; slug: string; count: number } | { kind: "skip"; slug: string } | { kind: "fail"; slug: string; reason: string };
type ComponentPhaseOutcome = | { kind: "ok"; slug: string; snapshots: ComponentSnapshot[] } | { kind: "skip"; slug: string } | { kind: "fail"; slug: string; reason: string };
type Triplet = { row: ExternalServiceRow; entry: StatusPageEntry; fetcher: StatusFetcher | null;};
function runStatusPhase( triplets: Triplet[], fetchedAt: number,): Effect.Effect<StatusPhaseOutcome[]> { return Effect.forEach( triplets, ({ entry, fetcher }) => { if (!fetcher) { return Effect.succeed<StatusPhaseOutcome>({ kind: "no-fetcher", slug: entry.id, }); } return fetcher.fetch(entry).pipe( Effect.map((result): StatusPhaseOutcome => ({ kind: "ok", snapshot: buildSnapshot({ entry, result, fetchedAt }), })), // Failure reporting is deferred: the detect step after this phase // either merges it into a detection story or reports it plain. Effect.catch((err: FetchError) => Effect.succeed<StatusPhaseOutcome>({ kind: "fail", slug: entry.id, reason: err.message, error: err, }), ), ); }, { concurrency: PHASE_CONCURRENCY }, );}
function runIncidentPhase( triplets: Triplet[], tickStartedAt: Date,): Effect.Effect<IncidentPhaseOutcome[]> { return Effect.forEach( triplets, ({ row, entry, fetcher }) => { if (!fetcher || !fetcher.fetchIncidents) { return Effect.succeed<IncidentPhaseOutcome>({ kind: "skip", slug: entry.id, }); } return fetcher.fetchIncidents(entry).pipe( Effect.flatMap((incidents) => Effect.tryPromise({ try: () => upsertExternalIncidentsForService({ ctx: { db }, externalServiceId: row.id, incidents: incidents.map(toUpsertInput), now: tickStartedAt, }), catch: (e) => new FetchError({ url: entry.status_page_url, fetcherName: fetcher.name, entryId: entry.id, cause: e instanceof Error ? e : new Error(String(e)), }), }).pipe( Effect.map((result): IncidentPhaseOutcome => ({ kind: "ok", slug: entry.id, count: result.upserted, })), ), ), Effect.catch((err: FetchError) => Effect.sync<IncidentPhaseOutcome>(() => { reportFetchFailure({ phase: "incidents", slug: entry.id, error: err, level: isBlocked(err) ? "warning" : undefined, }); return { kind: "fail", slug: entry.id, reason: err.message }; }), ), ); }, { concurrency: PHASE_CONCURRENCY }, );}
function runComponentPhase( triplets: Triplet[], tickStartedAt: Date, fetchedAt: number,): Effect.Effect<ComponentPhaseOutcome[]> { return Effect.forEach( triplets, ({ row, entry, fetcher }) => { if (!fetcher || !fetcher.fetchComponents) { return Effect.succeed<ComponentPhaseOutcome>({ kind: "skip", slug: entry.id, }); } return fetcher.fetchComponents(entry).pipe( Effect.flatMap((components) => Effect.tryPromise({ try: () => upsertExternalComponentsForService({ ctx: { db }, externalServiceId: row.id, components: components.map(toComponentUpsertInput), now: tickStartedAt, }), catch: (e) => new FetchError({ url: entry.status_page_url, fetcherName: fetcher.name, entryId: entry.id, cause: e instanceof Error ? e : new Error(String(e)), }), }).pipe( Effect.map((result): ComponentPhaseOutcome => { // History rows key on our PK, so the upstreamโPK map from the // upsert turns each normalized component into a snapshot. const byUpstream = new Map( components.map((c) => [c.upstreamComponentId, c]), ); const snapshots: ComponentSnapshot[] = []; for (const upserted of result.upserted) { const c = byUpstream.get(upserted.upstreamComponentId); if (!c) continue; snapshots.push({ component_id: String(upserted.id), external_service_id: row.id, indicator: c.severity, status: c.status, fetched_at: fetchedAt, }); } return { kind: "ok", slug: entry.id, snapshots }; }), ), ), Effect.catch((err: FetchError) => Effect.sync<ComponentPhaseOutcome>(() => { reportFetchFailure({ phase: "components", slug: entry.id, error: err, level: isBlocked(err) ? "warning" : undefined, }); return { kind: "fail", slug: entry.id, reason: err.message }; }), ), ); }, { concurrency: PHASE_CONCURRENCY }, );}
function summarizeStatus(outcomes: StatusPhaseOutcome[]): { counts: PhaseCounts; snapshots: Snapshot[];} { const snapshots: Snapshot[] = []; let successCount = 0; let failureCount = 0; let skippedCount = 0; for (const o of outcomes) { if (o.kind === "ok") { snapshots.push(o.snapshot); successCount++; } else if (o.kind === "no-fetcher") { skippedCount++; logger.warn("external-status status: no fetcher matches slug={slug}", { slug: o.slug, }); } else { failureCount++; logger.warn( "external-status status: fetch failed for slug={slug}: {reason}", { slug: o.slug, reason: o.reason, }, ); } } return { counts: { successCount, failureCount, skippedCount, total: outcomes.length, }, snapshots, };}
function summarizeIncidents(outcomes: IncidentPhaseOutcome[]): PhaseCounts { let successCount = 0; let failureCount = 0; let skippedCount = 0; for (const o of outcomes) { if (o.kind === "ok") { successCount++; } else if (o.kind === "skip") { skippedCount++; } else { failureCount++; logger.warn( "external-status incidents: failed for slug={slug}: {reason}", { slug: o.slug, reason: o.reason, }, ); } } return { successCount, failureCount, skippedCount, total: outcomes.length, };}
function summarizeComponents(outcomes: ComponentPhaseOutcome[]): { counts: PhaseCounts; snapshots: ComponentSnapshot[];} { const snapshots: ComponentSnapshot[] = []; let successCount = 0; let failureCount = 0; let skippedCount = 0; for (const o of outcomes) { if (o.kind === "ok") { successCount++; snapshots.push(...o.snapshots); } else if (o.kind === "skip") { skippedCount++; } else { failureCount++; logger.warn( "external-status components: failed for slug={slug}: {reason}", { slug: o.slug, reason: o.reason, }, ); } } return { counts: { successCount, failureCount, skippedCount, total: outcomes.length, }, snapshots, };}
function buildTriplets(services: ExternalServiceRow[]): Triplet[] { return services.map((row) => { const entry = toStatusPageEntry(row); const fetcher = fetchers.find((f) => f.canHandle(entry)) ?? null; return { row, entry, fetcher }; });}
const DETECT_CONCURRENCY = 3;
type DetectItem = { triplet: Triplet; error?: FetchError };
type DetectOutcome = { kind: | "applied" | "config-fixed" | "suggested" | "transient" | "none" | "skipped" | "error"; slug: string;};
type DetectCounts = { probed: number; applied: number; configFixed: number; suggested: number; failed: number;};
function collectDetectItems( triplets: Triplet[], statusOutcomes: StatusPhaseOutcome[], now: number,): DetectItem[] { const items: DetectItem[] = []; statusOutcomes.forEach((outcome, i) => { const triplet = triplets[i]; if (!triplet) return; if (outcome.kind === "no-fetcher") { if (shouldProbe(triplet.entry.id, now)) items.push({ triplet }); return; } if (outcome.kind !== "fail") return; if (isSuspicious(outcome.error) && shouldProbe(triplet.entry.id, now)) { items.push({ triplet, error: outcome.error }); } else { reportFetchFailure({ phase: "status", slug: outcome.slug, error: outcome.error, level: isBlocked(outcome.error) ? "warning" : undefined, }); } }); return items;}
function applyDetection(args: { triplet: Triplet; error?: FetchError; provider: ExternalServiceRow["provider"]; outcome: "applied" | "config-fixed"; evidence: string[]; tickStartedAt: Date;}): Effect.Effect<DetectOutcome> { const { triplet, error, provider, outcome, evidence, tickStartedAt } = args; const { row, entry } = triplet; return Effect.tryPromise({ try: () => applyDetectedProvider({ ctx: { db }, input: { serviceId: row.id, expected: { provider: row.provider, apiConfig: row.apiConfig ?? null, }, set: { provider, apiConfig: null }, }, now: tickStartedAt, }), catch: (e) => (e instanceof Error ? e : new Error(String(e))), }).pipe( Effect.map(({ updated }): DetectOutcome => { if (!updated) { logger.warn( "external-status detect: concurrent edit, write skipped for slug={slug}", { slug: entry.id }, ); return { kind: "skipped", slug: entry.id }; } reportDetectionStory({ slug: entry.id, currentProvider: row.provider, fetchError: error, outcome: outcome === "applied" ? { kind: "applied", provider } : { kind: "config-cleared" }, evidence, }); return { kind: outcome, slug: entry.id }; }), Effect.catch((e) => Effect.sync((): DetectOutcome => { logger.warn( "external-status detect: write failed for slug={slug}: {message}", { slug: entry.id, message: e.message }, ); reportDetectionWriteFailure({ slug: entry.id, error: e }); clearProbeStamp(entry.id); // The merged story never fired; keep the triggering failure visible. if (error) { reportFetchFailure({ phase: "status", slug: entry.id, error }); } return { kind: "error", slug: entry.id }; }), ), );}
function detectAndAct( item: DetectItem, tickStartedAt: Date,): Effect.Effect<DetectOutcome> { const { triplet, error } = item; const { row, entry } = triplet; return detectProvider({ statusPageUrl: entry.status_page_url, currentProvider: row.provider, entryId: entry.id, }).pipe( Effect.flatMap((result) => { const action = decideDetectionAction(result, row); switch (action.kind) { case "apply": return applyDetection({ triplet, error, provider: action.provider, outcome: "applied", evidence: action.evidence, tickStartedAt, }); case "clear-config": return applyDetection({ triplet, error, provider: row.provider, outcome: "config-fixed", evidence: action.evidence, tickStartedAt, }); case "suggest": return Effect.sync((): DetectOutcome => { reportDetectionStory({ slug: entry.id, currentProvider: row.provider, fetchError: error, outcome: { kind: "suggest", suggestion: action.suggestion }, evidence: action.evidence, }); return { kind: "suggested", slug: entry.id }; }); case "noop": return Effect.sync((): DetectOutcome => { if (action.reason === "no-evidence") { // Valid JSON our schema rejects means our schema is stale, not // that the page moved. reportDetectionStory({ slug: entry.id, currentProvider: row.provider, fetchError: error, outcome: error?.kind === "schema" ? { kind: "schema-mismatch" } : { kind: "none" }, evidence: action.evidence, }); return { kind: "none", slug: entry.id }; } if (error) { reportFetchFailure({ phase: "status", slug: entry.id, error }); } return { kind: "transient", slug: entry.id }; }); } }), );}
function runDetectPhase( items: DetectItem[], tickStartedAt: Date,): Effect.Effect<DetectOutcome[]> { return Effect.forEach(items, (item) => detectAndAct(item, tickStartedAt), { concurrency: DETECT_CONCURRENCY, });}
function summarizeDetect(outcomes: DetectOutcome[]): DetectCounts { let applied = 0; let configFixed = 0; let suggested = 0; let failed = 0; for (const o of outcomes) { if (o.kind === "applied") applied++; else if (o.kind === "config-fixed") configFixed++; else if (o.kind === "suggested") suggested++; else if (o.kind === "error") failed++; } return { probed: outcomes.length, applied, configFixed, suggested, failed };}
export async function runExternalStatusTick(): Promise<{ status: PhaseCounts; incidents: PhaseCounts; components: PhaseCounts; detect: DetectCounts;}> { const services = await listExternalServices({ ctx: { db } });
const triplets = buildTriplets(services); const tickStartedAt = new Date();
// The three phases are intentionally independent: each hits a different // upstream endpoint and store, so a failed status fetch must not suppress a // service's components (or vice versa). A tick can therefore persist // component history for a service whose status snapshot failed that tick. const [statusOutcomes, incidentOutcomes, componentOutcomes] = await Effect.runPromise( Effect.all( [ runStatusPhase(triplets, tickStartedAt.getTime()), runIncidentPhase(triplets, tickStartedAt), runComponentPhase(triplets, tickStartedAt, tickStartedAt.getTime()), ], { concurrency: 50 }, ), );
const status = summarizeStatus(statusOutcomes); const incidents = summarizeIncidents(incidentOutcomes); const components = summarizeComponents(componentOutcomes);
if (status.snapshots.length > 0) { await tb.publishExternalStatus(status.snapshots); } if (components.snapshots.length > 0) { await tb.publishExternalStatusComponent(components.snapshots); }
// After the publishes: a slow probe fleet must not delay or drop the // snapshots the tick already fetched. const detectItems = collectDetectItems(triplets, statusOutcomes, Date.now()); const detectOutcomes = detectItems.length > 0 ? await Effect.runPromise(runDetectPhase(detectItems, tickStartedAt)) : []; const detect = summarizeDetect(detectOutcomes);
return { status: status.counts, incidents, components: components.counts, detect, };}
export async function handleExternalStatusCron(c: Context) { const { cronCompleted, cronFailed } = runSentryCron("external-status");
// Background chain: must not capture `c` or anything derived from it // (e.g. via getSentry(c)). The handler returns 200 before this resolves, and // a captured per-request Sentry hub stays pinned across retries โ see // apps/workflows/plan.md. void Effect.runPromise( Effect.tryPromise({ try: () => runExternalStatusTick(), catch: (e) => new Error( `external-status tick failed: ${e instanceof Error ? e.message : String(e)}`, ), }).pipe( Effect.tap((res) => Effect.sync(() => { logger.info( "external-status tick complete: status={statusOk}/{statusTotal} ({statusFail} failures, {statusSkip} skipped), incidents={incOk}/{incTotal} ({incFail} failures, {incSkip} skipped), components={compOk}/{compTotal} ({compFail} failures, {compSkip} skipped), detect={detProbed} probed ({detApplied} applied, {detCfg} config-fixed, {detSuggested} suggested, {detFailed} failed)", { detProbed: res.detect.probed, detApplied: res.detect.applied, detCfg: res.detect.configFixed, detSuggested: res.detect.suggested, detFailed: res.detect.failed, statusOk: res.status.successCount, statusTotal: res.status.total, statusFail: res.status.failureCount, statusSkip: res.status.skippedCount, incOk: res.incidents.successCount, incTotal: res.incidents.total, incFail: res.incidents.failureCount, incSkip: res.incidents.skippedCount, compOk: res.components.successCount, compTotal: res.components.total, compFail: res.components.failureCount, compSkip: res.components.skippedCount, }, ); void cronCompleted(); }), ), Effect.catch((e) => Effect.sync(() => { logger.error("external-status tick errored: {message}", { message: e.message, }); void reportBackgroundError(e.message); void cronFailed(); }), ), ), );
return c.json({ success: true }, 200);}