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.7 kB ยท 251 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252import { and, count, db as defaultDb, eq } from "@openstatus/db";import { page, pageComponent } from "@openstatus/db/src/schema";import type { ImportSummary } from "@openstatus/importers";
import { requireScope } from "../auth";import type { ServiceContext } from "../context";import { NotFoundError, ValidationError } from "../errors";import { addLimitWarnings } from "./limits";import { type PhaseContext, writeComponentGroupsPhase, writeComponentsPhase, writeIncidentsPhase, writeMaintenancesPhase, writeMonitorsPhase, writePagePhase, writeSubscribersPhase,} from "./phase-writers";import { buildProviderConfig, createProvider } from "./provider";import { RunImportInput } from "./schemas";
/** * Execute a real import: fetches from the provider, validates, applies * limit warnings, then walks each phase in sequence writing to the db. * * The orchestrator is deliberately *not* wrapped in a single * `withTransaction`: imports can span many minutes and hold locks across * dozens of writes, and the existing UX is phase-level recovery โ a * failing phase aborts subsequent phases but preserves earlier phases' * writes. * * Audit emission: each phase writer emits per-resource rows * (`page.create`, `monitor.create`, etc.) for every resource it * actually creates, matching what the domain services would emit for * normal CRUD. Skipped rows have their original create audit already; * failed rows have nothing to attribute. No rollup row โ the action * union only allows create/update/delete verbs per entity. */export async function runImport(args: { ctx: ServiceContext; input: RunImportInput;}): Promise<ImportSummary> { const { ctx } = args; requireScope(ctx, "write"); const input = RunImportInput.parse(args.input); const tx = ctx.db ?? defaultDb;
// Verify pageId belongs to the workspace before doing any provider work. // Was previously duplicated at the router layer โ owning this in the // service means all callers (tRPC / Slack / future) get the same check. if (input.pageId) { const existing = await tx .select({ id: page.id }) .from(page) .where( and(eq(page.id, input.pageId), eq(page.workspaceId, ctx.workspace.id)), ) .get();
if (!existing) throw new NotFoundError("page", input.pageId); }
const provider = createProvider(input.provider); const providerConfig = buildProviderConfig({ provider: input.provider, apiKey: input.apiKey, statuspagePageId: input.statuspagePageId, betterstackStatusPageId: input.betterstackStatusPageId, instatusPageId: input.instatusPageId, checklyAccountId: input.checklyAccountId, checklyStatusPageId: input.checklyStatusPageId, workspaceId: ctx.workspace.id, pageId: input.pageId, });
const validation = await provider.validate(providerConfig); if (!validation.valid) { throw new ValidationError( `Provider validation failed: ${validation.error ?? "unknown error"}`, ); }
const summary = await provider.run(providerConfig);
await addLimitWarnings(summary, { limits: ctx.workspace.limits, workspaceId: ctx.workspace.id, pageId: input.pageId, db: tx, options: input.options, });
const idMaps = { groups: new Map<string, number>(), components: new Map<string, number>(), monitors: new Map<string, number>(), };
const pc: PhaseContext = { ctx, tx, provider: input.provider };
let targetPageId = input.pageId; let phaseAborted = false;
for (const phase of summary.phases) { if (phaseAborted) { phase.status = "skipped"; continue; }
try { switch (phase.phase) { case "monitors": if (input.options?.includeMonitors !== false) { await writeMonitorsPhase(pc, phase, idMaps.monitors); } else { phase.status = "skipped"; } break; case "page": targetPageId = await writePagePhase(pc, phase, input.pageId); break; case "componentGroups": if (targetPageId && input.options?.includeComponents !== false) { await writeComponentGroupsPhase( pc, phase, targetPageId, idMaps.groups, ); } else { // Fall-through skip: either `includeComponents === false` // (user opt-out) or `targetPageId` is missing (page phase // produced nothing). Either way this phase has nothing to do. phase.status = "skipped"; } break; case "components": if (targetPageId && input.options?.includeComponents !== false) { // Workspace-wide count โ `page-components` is the plan cap // across every page in the workspace (see // `page-component/update-order`). Scoping to `targetPageId` // alone would let an import into an empty page push the // workspace past the cap because components on other // pages go uncounted. const [compCount] = await tx .select({ count: count() }) .from(pageComponent) .where(eq(pageComponent.workspaceId, ctx.workspace.id)); const maxComponents = ctx.workspace.limits["page-components"]; const remaining = maxComponents - (compCount?.count ?? 0); if (remaining <= 0) { // Stamp each resource with a skip reason to mirror the // `writeMonitorsPhase` pattern โ otherwise users see a // failed phase with no explanation and stale resource // statuses. for (const r of phase.resources) { r.status = "skipped"; r.error = `Skipped: component limit reached (${maxComponents})`; } phase.status = "skipped"; break; } if (phase.resources.length > remaining) { // Trim the overflow: mark skipped with a clear error string // and keep them in `resources` so the summary still reports // them. const skipped = phase.resources.splice(remaining); for (const r of skipped) { r.status = "skipped"; r.error = `Skipped: would exceed component limit (${maxComponents})`; } phase.resources.push(...skipped); } await writeComponentsPhase( pc, phase, targetPageId, idMaps.groups, idMaps.components, idMaps.monitors, ); } else { phase.status = "skipped"; } break; case "incidents": if (targetPageId && input.options?.includeStatusReports !== false) { await writeIncidentsPhase( pc, phase, targetPageId, idMaps.components, ); } else { phase.status = "skipped"; } break; case "maintenances": if (targetPageId && input.options?.includeStatusReports !== false) { await writeMaintenancesPhase( pc, phase, targetPageId, idMaps.components, ); } else { phase.status = "skipped"; } break; case "subscribers": if (targetPageId && input.options?.includeSubscribers) { if (!ctx.workspace.limits["status-subscribers"]) { phase.status = "skipped"; break; } await writeSubscribersPhase( pc, phase, targetPageId, idMaps.components, ); } else { phase.status = "skipped"; } break; } } catch (err) { const msg = err instanceof Error ? err.message : String(err); summary.errors.push(`Phase "${phase.phase}" failed: ${msg}`); phase.status = "failed"; phaseAborted = true; } }
const hasFailures = summary.phases.some((p) => p.status === "failed"); const hasPartial = summary.phases.some((p) => p.status === "partial"); const allSkippedOrCompleted = summary.phases.every( (p) => p.status === "completed" || p.status === "skipped", );
summary.status = hasFailures ? "failed" : hasPartial ? "partial" : allSkippedOrCompleted ? "completed" : "partial"; summary.completedAt = new Date();
return summary;}