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.
15 kB ยท 555 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556import { CloudTasksClient } from "@google-cloud/tasks";import type { google } from "@google-cloud/tasks/build/protos";import { and, db, desc, eq, isNull, lte, max, or, schema,} from "@openstatus/db";import { selectWorkspaceSchema, user } from "@openstatus/db/src/schema";import { monitorDeactivationEmail, monitorPausedEmail,} from "@openstatus/emails";import { type EmailHtml, sendBatchEmailHtml,} from "@openstatus/emails/src/send";import { bulkUpdateMonitors } from "@openstatus/services/monitor";import { Redis } from "@openstatus/upstash";import { RateLimiter } from "limiter";import { z } from "zod";
import { env } from "../env";
const redis = Redis.fromEnv();
const PAUSE_FROM = "Thibault from openstatus <thibault@notifications.openstatus.dev>";const PAUSE_REPLY_TO = "thibault@openstatus.dev";
// Check if GCP is properly configured (not empty and not placeholder values)const isGcpConfigured = () => { const gcpProjectId = env().GCP_PROJECT_ID; const gcpLocation = env().GCP_LOCATION; const gcpClientEmail = env().GCP_CLIENT_EMAIL; const gcpPrivateKey = env().GCP_PRIVATE_KEY;
return ( gcpProjectId && gcpLocation && gcpClientEmail && gcpPrivateKey && gcpProjectId !== "" && gcpLocation !== "" && gcpClientEmail !== "" && gcpPrivateKey !== "" && gcpProjectId !== "your-value" && gcpLocation !== "your-value" && gcpClientEmail !== "your-value" && gcpPrivateKey !== "your-value" );};
// Only initialize Cloud Tasks client if GCP is properly configuredlet client: CloudTasksClient | null = null;let parent: string | null = null;
if (isGcpConfigured()) { client = new CloudTasksClient({ projectId: env().GCP_PROJECT_ID, fallback: "rest", credentials: { client_email: env().GCP_CLIENT_EMAIL, private_key: env().GCP_PRIVATE_KEY.replaceAll("\\n", "\n"), }, });
parent = client.queuePath( env().GCP_PROJECT_ID, env().GCP_LOCATION, "workflow", );}
const limiter = new RateLimiter({ tokensPerInterval: 15, interval: "second" });
export async function LaunchMonitorWorkflow() { // Expires is one month after last connection, so if we want to reach people who connected 3 months ago we need to check for people with expires 2 months ago const twoMonthAgo = new Date().setMonth(new Date().getMonth() - 2);
const date = new Date(twoMonthAgo); // User without session const userWithoutSession = db .select({ userId: schema.user.id, email: schema.user.email, updatedAt: schema.user.updatedAt, }) .from(schema.user) .leftJoin(schema.session, eq(schema.session.userId, schema.user.id)) .where(isNull(schema.session.userId)) .as("query"); // Only free users monitors are paused // We don't need to handle multi users per workspace because free workspaces only have one user // Only free users monitors are paused
const u1 = await db .select({ userId: userWithoutSession.userId, email: userWithoutSession.email, workspaceId: schema.workspace.id, }) .from(userWithoutSession) .innerJoin( schema.usersToWorkspaces, eq(userWithoutSession.userId, schema.usersToWorkspaces.userId), ) .innerJoin( schema.workspace, eq(schema.usersToWorkspaces.workspaceId, schema.workspace.id), ) .where( and( or( lte(userWithoutSession.updatedAt, date), isNull(userWithoutSession.updatedAt), ), or(isNull(schema.workspace.plan), eq(schema.workspace.plan, "free")), ), );
console.log(`Found ${u1.length} users without session to start the workflow`); const maxSessionPerUser = db .select({ userId: schema.user.id, email: schema.user.email, lastConnection: max(schema.session.expires).as("lastConnection"), }) .from(schema.user) .innerJoin(schema.session, eq(schema.session.userId, schema.user.id)) .groupBy(schema.user.id) .as("maxSessionPerUser"); // Only free users monitors are paused // We don't need to handle multi users per workspace because free workspaces only have one user // Only free users monitors are paused
const u = await db .select({ userId: maxSessionPerUser.userId, email: maxSessionPerUser.email, workspaceId: schema.workspace.id, }) .from(maxSessionPerUser) .innerJoin( schema.usersToWorkspaces, eq(maxSessionPerUser.userId, schema.usersToWorkspaces.userId), ) .innerJoin( schema.workspace, eq(schema.usersToWorkspaces.workspaceId, schema.workspace.id), ) .where( and( lte(maxSessionPerUser.lastConnection, date), or(isNull(schema.workspace.plan), eq(schema.workspace.plan, "free")), ), ); const usersMap = new Map<number, (typeof u)[number]>(); for (const entry of [...u, ...u1]) { usersMap.set(entry.userId, entry); } const users = Array.from(usersMap.values()); const duplicatesRemoved = u.length + u1.length - users.length; if (duplicatesRemoved > 0) { console.log(`Removed ${duplicatesRemoved} duplicate users`); }
const allResult = [];
for (const user of users) { await limiter.removeTokens(1); const workflow = workflowInit({ user }); allResult.push(workflow); }
const allRequests = await Promise.allSettled(allResult);
const success = allRequests.filter((r) => r.status === "fulfilled").length; const failed = allRequests.filter((r) => r.status === "rejected").length;
console.log( `End cron with ${allResult.length} jobs with ${success} success and ${failed} failed`, );}
async function workflowInit({ user,}: { user: { userId: number; email: string | null; workspaceId: number; };}) { console.log(`Starting workflow for ${user.userId}`); const isMember = await redis.exists(`workflow:user:${user.userId}`); if (isMember) { console.log(`user workflow already started for ${user.userId}`); return; } // check if user has some running monitors const nbRunningMonitor = await db.$count( schema.monitor, and( eq(schema.monitor.workspaceId, user.workspaceId), eq(schema.monitor.active, true), isNull(schema.monitor.deletedAt), ), ); if (nbRunningMonitor === 0) { console.log(`user has no running monitors for ${user.userId}`); return; } const initialRun = new Date().getTime();
// Only create GCP Cloud Tasks if GCP is configured if (!client || !parent) { console.log( `Skipping workflow for user ${user.userId} - GCP not configured`, ); return; }
await CreateTask({ parent, client: client, step: "14days", userId: user.userId, initialRun, });
await redis.set(`workflow:user:${user.userId}`, initialRun, { ex: 30 * 86400, }); console.log(`user workflow started for ${user.userId}`);}
export async function Step14Days(userId: number, workFlowRunTimestamp: number) { const hasConnected = await hasUserLoggedIn({ userId, date: new Date(workFlowRunTimestamp), });
if (hasConnected) { await redis.del(`workflow:user:${userId}`); return; }
const user = await getUser(userId);
if (user.email) { const email = await monitorDeactivationEmail({ deactivateAt: new Date(new Date().setDate(new Date().getDate() + 14)), ...(await getPauseContext(userId)), }); await sendWorkflowEmail({ userId, step: "14days", initialRun: workFlowRunTimestamp, email: { to: user.email, from: PAUSE_FROM, reply_to: PAUSE_REPLY_TO, ...email, }, }); }
// Only create GCP Cloud Tasks if GCP is configured if (client && parent) { await CreateTask({ parent, client: client, step: "3days", userId: user.id, initialRun: workFlowRunTimestamp, }); } else { console.log( `Skipping Cloud Tasks creation for user ${user.id} - GCP not configured`, ); }}
export async function Step3Days(userId: number, workFlowRunTimestamp: number) { const hasConnected = await hasUserLoggedIn({ userId, date: new Date(workFlowRunTimestamp), });
if (hasConnected) { await redis.del(`workflow:user:${userId}`); return; }
const user = await getUser(userId);
if (user.email) { const email = await monitorDeactivationEmail({ deactivateAt: new Date(new Date().setDate(new Date().getDate() + 3)), ...(await getPauseContext(userId)), }); await sendWorkflowEmail({ userId, step: "3days", initialRun: workFlowRunTimestamp, email: { to: user.email, from: PAUSE_FROM, reply_to: PAUSE_REPLY_TO, ...email, }, }); }
// Only create GCP Cloud Tasks if GCP is configured if (client && parent) { await CreateTask({ client, parent, step: "paused", userId, initialRun: workFlowRunTimestamp, }); } else { console.log( `Skipping Cloud Tasks creation for user ${userId} - GCP not configured`, ); }}
export async function StepPaused(userId: number, workFlowRunTimestamp: number) { try { const hasConnected = await hasUserLoggedIn({ userId, date: new Date(workFlowRunTimestamp), });
if (hasConnected) { return; }
const workspace = await getUserWorkspace(userId); let pausedCount = 0; if (workspace) { const activeMonitors = await db .select({ id: schema.monitor.id }) .from(schema.monitor) .where( and( eq(schema.monitor.workspaceId, workspace.id), eq(schema.monitor.active, true), isNull(schema.monitor.deletedAt), ), ) .all(); pausedCount = activeMonitors.length; if (activeMonitors.length > 0) { await bulkUpdateMonitors({ ctx: { workspace, actor: { type: "system", job: "monitor-auto-pause" }, }, input: { ids: activeMonitors.map((m) => m.id), active: false }, }); } }
const currentUser = await getUser(userId); if (currentUser.email) { const email = await monitorPausedEmail({ monitorCount: pausedCount || undefined, workspaceSlug: workspace?.slug, }); await sendWorkflowEmail({ userId, step: "paused", initialRun: workFlowRunTimestamp, email: { to: currentUser.email, from: PAUSE_FROM, reply_to: PAUSE_REPLY_TO, ...email, }, }); } } finally { await redis.del(`workflow:user:${userId}`); }}
async function hasUserLoggedIn({ userId, date,}: { userId: number; date: Date;}) { const userResult = await db .select({ lastSession: schema.session.expires }) .from(schema.session) .where(eq(schema.session.userId, userId)) .orderBy(desc(schema.session.expires));
if (userResult.length === 0) { return false; } const user = userResult[0]; if (user.lastSession === null) { return false; } return user.lastSession > date;}
async function CreateTask({ parent, client, step, userId, initialRun,}: { parent: string; client: CloudTasksClient; step: z.infer<typeof workflowStepSchema>; userId: number; initialRun: number;}) { const url = `https://openstatus-workflows.fly.dev/cron/monitors/${step}?userId=${userId}&initialRun=${initialRun}`; const timestamp = getScheduledTime(step); const taskName = `${parent}/tasks/workflow-${userId}-${step}-${initialRun}`; const newTask: google.cloud.tasks.v2beta3.ITask = { name: taskName, httpRequest: { headers: { "Content-Type": "application/json", Authorization: `${env().CRON_SECRET}`, }, httpMethod: "GET", url, }, scheduleTime: { seconds: timestamp, }, };
const request = { parent: parent, task: newTask }; try { return await client.createTask(request); } catch (e) { if (e instanceof Error && "code" in e && e.code === 6) { console.log( `Task already exists for user ${userId} step ${step}, skipping`, ); return; } throw e; }}
function getScheduledTime(step: z.infer<typeof workflowStepSchema>) { switch (step) { case "14days": // let's triger it now return new Date().getTime() / 1000; case "3days": // it's 11 days after the 14 days return new Date().setDate(new Date().getDate() + 11) / 1000; case "paused": // it's 3 days after the 3 days step return new Date().setDate(new Date().getDate() + 3) / 1000; default: throw new Error("Invalid step"); }}
export const workflowStep = ["14days", "3days", "paused"] as const;export const workflowStepSchema = z.enum(workflowStep);
async function getUser(userId: number) { const currentUser = await db .select() .from(user) .where(eq(schema.user.id, userId)) .get();
if (!currentUser) { throw new Error("User not found"); } return currentUser;}
async function sendWorkflowEmail({ userId, step, initialRun, email,}: { userId: number; step: string; initialRun: number; email: EmailHtml;}) { const key = `workflow:email:${userId}:${step}:${initialRun}`; const alreadySent = await redis.exists(key); if (alreadySent) { console.log(`Email already sent for user ${userId} step ${step}, skipping`); return; } await sendBatchEmailHtml([email]); await redis.set(key, 1, { ex: 30 * 86400 });}
async function getPauseContext(userId: number) { const workspace = await getUserWorkspace(userId); if (!workspace) return {}; const monitorCount = await db.$count( schema.monitor, and( eq(schema.monitor.workspaceId, workspace.id), eq(schema.monitor.active, true), isNull(schema.monitor.deletedAt), ), ); return { monitorCount: monitorCount || undefined, workspaceSlug: workspace.slug, };}
async function getUserWorkspace(userId: number) { const row = await db .select({ workspace: schema.workspace }) .from(schema.user) .innerJoin( schema.usersToWorkspaces, eq(schema.user.id, schema.usersToWorkspaces.userId), ) .innerJoin( schema.workspace, eq(schema.usersToWorkspaces.workspaceId, schema.workspace.id), ) .where( and( eq(schema.user.id, userId), or(isNull(schema.workspace.plan), eq(schema.workspace.plan, "free")), ), ) .get(); return row ? selectWorkspaceSchema.parse(row.workspace) : undefined;}