diff --git a/AGENTS.md b/AGENTS.md index e4024ff..a732d09 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -178,7 +178,7 @@ Requires secrets: `VERCEL_TOKEN`, `VERCEL_ORG_ID`, `VERCEL_PROJECT_ID` | `DATABASE_URL` | `postgresql://dequel:dequel@localhost:5432/dequel` | PostgreSQL connection string | | `WORKSPACE_ROOT` | `./workspace` | Build staging | | `CADDY_ROUTES_DIR` | `./infra/caddy/routes` | Caddy route output | -| `CADDY_BASE_DOMAIN` | `localhost` | Base domain for deployment subdomains. Set to a real domain (e.g. `example.com`) for Let's Encrypt auto-SSL. | +| `CADDY_BASE_DOMAIN` | `localhost` | Base domain for deployment subdomains. Set to a real domain (e.g. `example.com`) for Let's Encrypt auto-SSL. Public links (e.g. failure email logs) derive their base URL from this. | | `CADDY_EMAIL` | _(empty)_ | Email for Let's Encrypt SSL certificate notifications | | `DOCKER_NETWORK` | `dequel_net` | Docker network for deployments | | `BUILDKIT_HOST` | `tcp://buildkit:1234` | Buildkit daemon | diff --git a/apps/api/src/agents/job-channel.ts b/apps/api/src/agents/job-channel.ts index 8c0e4af..c64bb91 100644 --- a/apps/api/src/agents/job-channel.ts +++ b/apps/api/src/agents/job-channel.ts @@ -9,6 +9,7 @@ import { leaseNextAgentJob, listCancelledJobIds, listDeployments, + recordDeploymentFailure, updateAgentHeartbeat, updateDeploymentCommitSha, updateDeploymentStatus, @@ -119,8 +120,10 @@ export const processAgentJobUpdate = async ( } } } else { - await updateDeploymentStatus(deploymentId, "failed", { - failureReason: update.error || "Remote agent deployment failed", + await recordDeploymentFailure({ + deploymentId, + reason: update.error || "Remote agent deployment failed", + source: "job-channel", }); await appendLog(deploymentId, "system", `Remote deployment failed: ${update.error || "Unknown agent error"}`); } diff --git a/apps/api/src/api/alerts/index.ts b/apps/api/src/api/alerts/index.ts index a1bc0c3..92e38be 100644 --- a/apps/api/src/api/alerts/index.ts +++ b/apps/api/src/api/alerts/index.ts @@ -1,7 +1,11 @@ import { Elysia } from "elysia"; import { createAlert, deleteAlert, listAlerts, updateAlertEnabled } from "../../db/repo"; +import type { AlertChannel, AlertType } from "../../types"; import { created, fail, ok } from "../response"; +const ALERT_TYPES: ReadonlySet = new Set(["cpu", "memory", "downtime", "cert_expiry"]); +const ALERT_CHANNELS: ReadonlySet = new Set(["email", "slack", "webhook"]); + export const alertsRoutes = new Elysia() .get("/projects/:id/alerts", async ({ params }) => ok(await listAlerts(params.id))) .post("/projects/:id/alerts", async ({ params, body, set }: any) => { @@ -9,6 +13,26 @@ export const alertsRoutes = new Elysia() set.status = 400; return fail("type and channel are required"); } + if (!ALERT_TYPES.has(String(body.type))) { + set.status = 400; + return fail(`type must be one of: ${[...ALERT_TYPES].join(", ")}`); + } + if (!ALERT_CHANNELS.has(String(body.channel))) { + set.status = 400; + return fail(`channel must be one of: ${[...ALERT_CHANNELS].join(", ")}`); + } + if (body.channel !== "email") { + let valid = false; + try { + valid = ["http:", "https:"].includes(new URL(String(body.destination ?? "")).protocol); + } catch { + valid = false; + } + if (!valid) { + set.status = 400; + return fail("destination must be an http(s) URL for this channel"); + } + } return created( await createAlert({ projectId: params.id, diff --git a/apps/api/src/api/backups/index.ts b/apps/api/src/api/backups/index.ts index 5bee6d6..b052cb0 100644 --- a/apps/api/src/api/backups/index.ts +++ b/apps/api/src/api/backups/index.ts @@ -1,7 +1,7 @@ import { Elysia } from "elysia"; import { BackupOrchestrator } from "../../backup/orchestrator"; -import { S3_BACKUP_PREFIX } from "../../backup/types"; import type { BackupTarget, StorageConfig } from "../../backup/types"; +import { S3_BACKUP_PREFIX } from "../../backup/types"; import { deleteBackupRecord, getBackupRecord, listBackupRecords } from "../../db/repo/backups"; import { getDatabaseById } from "../../db/repo/databases"; import { getBackupStorageSettings } from "../../db/repo/settings"; diff --git a/apps/api/src/api/deployments/index.ts b/apps/api/src/api/deployments/index.ts index 8146917..138e5c8 100644 --- a/apps/api/src/api/deployments/index.ts +++ b/apps/api/src/api/deployments/index.ts @@ -10,9 +10,11 @@ import { getProjectById, getServerById, listDeployments, + recordDeploymentFailure, } from "../../db/repo"; import { executorFor } from "../../executors/dispatch"; import { orchestrator } from "../../orchestrator"; +import { summarizeDeploymentError } from "../../orchestrator/deployment-errors"; import { logBus } from "../../orchestrator/log-bus"; import { config } from "../../utils/config"; import { isPrivateGitUrl } from "../../utils/validate"; @@ -31,8 +33,13 @@ const dispatchDeployment = async ( if (deployment.sourceType !== "git") throw new Error("Remote servers currently support Git deployments only"); const executor = executorFor(server.mode); if (server.mode === "ssh") { - void executor.deploy({ deployment, project, server }).catch((error) => { + void executor.deploy({ deployment, project, server }).catch(async (error) => { console.error(`[SSH Executor] Deployment ${deployment.id} failed:`, error); + await recordDeploymentFailure({ + deploymentId: deployment.id, + reason: summarizeDeploymentError(error), + source: "dispatch", + }).catch((e) => console.error(`[SSH Executor] Failed to record failure for ${deployment.id}:`, e)); }); return; } diff --git a/apps/api/src/api/index.ts b/apps/api/src/api/index.ts index b779741..c83b4ef 100644 --- a/apps/api/src/api/index.ts +++ b/apps/api/src/api/index.ts @@ -1,4 +1,5 @@ import { Elysia } from "elysia"; +import { fail } from "./response"; import { agentRoutes } from "./agents"; import { alertsRoutes } from "./alerts"; import { apiKeysRoutes } from "./api-keys"; @@ -11,6 +12,7 @@ import { envVarsRoutes } from "./env-vars"; import { githubRoutes } from "./github"; import { healthRoutes } from "./health"; import { projectsRoutes } from "./projects"; +import { projectStatusRoutes } from "./projects/status"; import { prometheusRoutes } from "./prometheus"; import { routesRoutes } from "./routes"; import { scalingRoutes } from "./scaling"; @@ -44,7 +46,7 @@ const authMiddleware = (app: Elysia) => const payload = await verifyAccessToken(match[1]); if (payload) return; set.status = 401; - return { error: "Invalid session" }; + return fail("Invalid session"); } const authHeader = request.headers.get("authorization"); @@ -55,22 +57,36 @@ const authMiddleware = (app: Elysia) => const key = await validateApiKey(token); if (key) return; set.status = 401; - return { error: "Invalid API key" }; + return fail("Invalid API key"); } } set.status = 401; - return { error: "Authentication required" }; + return fail("Authentication required"); }); +const INTERNAL_ERROR = + /Failed query:|params:|getaddrinfo|ECONNREFUSED|ETIMEDOUT|EAI_AGAIN|EHOSTUNREACH|timeout exceeded|node:internal|Cannot read propert|is not a function|is not a constructor|Unexpected token/i; + export const apiRoutes = new Elysia({ prefix: "/api", }) + .onError(({ error, set }) => { + const err = error as { status?: number; message?: string }; + set.status = typeof err?.status === "number" ? err.status : 500; + const message = err?.message ?? "Internal server error"; + if (INTERNAL_ERROR.test(message)) { + console.error("[API] Unhandled error:", error); + return fail("Internal server error"); + } + return fail(message); + }) .use(authRoutes) .use(authMiddleware) .use(agentRoutes) .use(healthRoutes) .use(projectsRoutes) + .use(projectStatusRoutes) .use(deploymentsRoutes) .use(envVarsRoutes) .use(sharedEnvVarsRoutes) diff --git a/apps/api/src/api/projects/index.ts b/apps/api/src/api/projects/index.ts index acc10ee..b8781f8 100644 --- a/apps/api/src/api/projects/index.ts +++ b/apps/api/src/api/projects/index.ts @@ -6,7 +6,6 @@ import { deleteProjectCascade, getProjectById, getServerById, - listDomains, listProjects, updateProject, } from "../../db/repo"; @@ -15,6 +14,7 @@ import { reloadCaddy, tryRun } from "../../orchestrator/runtime"; import { config } from "../../utils/config"; import { dockerBin } from "../../utils/docker-bin"; import { removeFromCaddyRoute } from "../../utils/domain-verifier"; +import { buildProjectRequestHostRegex, caddyRequestLogSelector } from "../../utils/loki"; import { isPort, isPrivateGitUrl, SERVICE_NAME_RE, validateComposeServices } from "../../utils/validate"; import { created, fail, ok } from "../response"; @@ -194,22 +194,8 @@ export const projectsRoutes = new Elysia() set.status = 404; return fail("Project not found"); } - const slugify = (s: string) => - s - .toLowerCase() - .replace(/[^a-z0-9-]+/g, "-") - .replace(/^-+|-+$/g, "") - .slice(0, 63); - const slug = slugify(project.name); - const domains = [`${slug}.${config.caddyBaseDomain}`]; - const projectDomains = await listDomains(id); - const verified = projectDomains.filter((d) => d.validationStatus === "verified"); - for (const d of verified) { - domains.push(d.domain); - } - - const regexEscaped = domains.map((d) => d.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\\\$&")).join("|"); - const queryStr = `{container="dequel-caddy-1"} | json | request_host =~ "^(${regexEscaped})$"`; + const hostRegex = await buildProjectRequestHostRegex(id); + const queryStr = caddyRequestLogSelector(hostRegex); const startParam = (queryParams as any)?.start; const endParam = (queryParams as any)?.end; @@ -274,23 +260,9 @@ export const projectsRoutes = new Elysia() set.status = 404; return fail("Project not found"); } - const slugify = (s: string) => - s - .toLowerCase() - .replace(/[^a-z0-9-]+/g, "-") - .replace(/^-+|-+$/g, "") - .slice(0, 63); - const slug = slugify(project.name); - const domains = [`${slug}.${config.caddyBaseDomain}`]; - const projectDomains = await listDomains(id); - const verified = projectDomains.filter((d) => d.validationStatus === "verified"); - for (const d of verified) { - domains.push(d.domain); - } + const hostRegex = await buildProjectRequestHostRegex(id); - const regexEscaped = domains.map((d) => d.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\\\$&")).join("|"); - - const query = `sum(count_over_time({container="dequel-caddy-1"} | json | request_host =~ "^(${regexEscaped})$" [5m]))`; + const query = `sum(count_over_time(${caddyRequestLogSelector(hostRegex)} [5m]))`; const end = Math.floor(Date.now() / 1000); const start = end - 6 * 60 * 60; @@ -325,22 +297,8 @@ export const projectsRoutes = new Elysia() return fail("Project not found"); } const encoder = new TextEncoder(); - const slugify = (s: string) => - s - .toLowerCase() - .replace(/[^a-z0-9-]+/g, "-") - .replace(/^-+|-+$/g, "") - .slice(0, 63); - const slug = slugify(project.name); - const domains = [`${slug}.${config.caddyBaseDomain}`]; - const projectDomains = await listDomains(id); - const verified = projectDomains.filter((d) => d.validationStatus === "verified"); - for (const d of verified) { - domains.push(d.domain); - } - - const regexEscaped = domains.map((d) => d.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\\\$&")).join("|"); - const query = `{container="dequel-caddy-1"} | json | request_host =~ "^(${regexEscaped})$"`; + const hostRegex = await buildProjectRequestHostRegex(id); + const query = caddyRequestLogSelector(hostRegex); let ws: WebSocket | null = null; let closed = false; diff --git a/apps/api/src/api/projects/status.ts b/apps/api/src/api/projects/status.ts new file mode 100644 index 0000000..eedf904 --- /dev/null +++ b/apps/api/src/api/projects/status.ts @@ -0,0 +1,29 @@ +import { Elysia } from "elysia"; +import { getProjectById } from "../../db/repo"; +import { getProjectStatus } from "../../monitoring/project-status"; +import { fail, ok } from "../response"; + +const DEFAULT_WINDOW_SECONDS = 3600; +const MIN_WINDOW_SECONDS = 60; +const MAX_WINDOW_SECONDS = 604800; + +const clampWindow = (raw: unknown): number => { + const value = Number(raw); + if (!Number.isFinite(value) || value <= 0) return DEFAULT_WINDOW_SECONDS; + return Math.min(MAX_WINDOW_SECONDS, Math.max(MIN_WINDOW_SECONDS, Math.floor(value))); +}; + +export const projectStatusRoutes = new Elysia().get("/projects/:id/status", async ({ params: { id }, query, set }) => { + const project = await getProjectById(id); + if (!project) { + set.status = 404; + return fail("Project not found"); + } + try { + return ok(await getProjectStatus(id, clampWindow((query as any)?.window))); + } catch (err) { + console.error(`[Status] Failed to build status for project ${id}:`, err); + set.status = 500; + return fail("Failed to build project status"); + } +}); diff --git a/apps/api/src/api/settings/index.ts b/apps/api/src/api/settings/index.ts index 61c5c81..79895c0 100644 --- a/apps/api/src/api/settings/index.ts +++ b/apps/api/src/api/settings/index.ts @@ -9,6 +9,7 @@ import { upsertBackupStorageSettings, upsertSmtpSettings, } from "../../db/repo"; +import { buildSmtpTestEmail } from "../../monitoring/templates"; import { failoverState } from "../../orchestrator/failover"; import { rerenderAllIngressRoutes } from "../../orchestrator/ingress-sync"; import { fail, ok } from "../response"; @@ -87,10 +88,12 @@ export const settingsRoutes = new Elysia({ prefix: "/settings" }) secure: settings.port === 465, auth: settings.user && settings.pass ? { user: settings.user, pass: settings.pass } : undefined, }); + const { subject, html } = buildSmtpTestEmail(); await transporter.sendMail({ from: settings.fromAddress, to: settings.fromAddress, - subject: "[Dequel] SMTP Test Email", + subject, + html, text: "This is a test email from Dequel. Your SMTP settings are working correctly.", }); return ok(null, "Test email sent"); diff --git a/apps/api/src/api/shared-env-vars/index.ts b/apps/api/src/api/shared-env-vars/index.ts index 9aa2095..d480956 100644 --- a/apps/api/src/api/shared-env-vars/index.ts +++ b/apps/api/src/api/shared-env-vars/index.ts @@ -1,7 +1,7 @@ -import { Elysia } from "elysia"; import { and, eq } from "drizzle-orm"; +import { Elysia } from "elysia"; import { getDb } from "../../db/db-provider"; -import { projectSharedEnvLinks } from "../../db/schema"; +import { getRowsAffected } from "../../db/repo/helpers"; import { createSharedEnvVar, deleteSharedEnvVar, @@ -12,7 +12,7 @@ import { listSharedEnvVars, updateSharedEnvVar, } from "../../db/repo/shared-env-vars"; -import { getRowsAffected } from "../../db/repo/helpers"; +import { projectSharedEnvLinks } from "../../db/schema"; import { fail, ok } from "../response"; export const sharedEnvVarsRoutes = new Elysia() @@ -89,12 +89,7 @@ export const sharedEnvLinksRoutes = new Elysia() getRowsAffected( await db .delete(projectSharedEnvLinks) - .where( - and( - eq(projectSharedEnvLinks.id, params.linkId), - eq(projectSharedEnvLinks.projectId, params.id), - ), - ) + .where(and(eq(projectSharedEnvLinks.id, params.linkId), eq(projectSharedEnvLinks.projectId, params.id))) .execute(), ) > 0; if (!removed) { diff --git a/apps/api/src/backup/__tests__/cron.test.ts b/apps/api/src/backup/__tests__/cron.test.ts new file mode 100644 index 0000000..1a721b3 --- /dev/null +++ b/apps/api/src/backup/__tests__/cron.test.ts @@ -0,0 +1,76 @@ +import { describe, expect, test } from "bun:test"; +import { matchesCron } from "../scheduler"; + +const at = (min: number, hour = 0, day = 1, month = 1) => new Date(2026, month - 1, day, hour, min); + +describe("matchesCron", () => { + test("star matches every value", () => { + for (const m of [0, 7, 59]) expect(matchesCron("* * * * *", at(m))).toBe(true); + }); + + test("step without base: */5 matches multiples of 5 only", () => { + expect(matchesCron("*/5 * * * *", at(0))).toBe(true); + expect(matchesCron("*/5 * * * *", at(5))).toBe(true); + expect(matchesCron("*/5 * * * *", at(55))).toBe(true); + expect(matchesCron("*/5 * * * *", at(3))).toBe(false); + expect(matchesCron("*/5 * * * *", at(7))).toBe(false); + }); + + test("range with step: 0-30/5 matches 0,5,...,30 but not 35", () => { + expect(matchesCron("0-30/5 * * * *", at(0))).toBe(true); + expect(matchesCron("0-30/5 * * * *", at(30))).toBe(true); + expect(matchesCron("0-30/5 * * * *", at(35))).toBe(false); + expect(matchesCron("0-30/5 * * * *", at(31))).toBe(false); + }); + + test("step with base: 5/15 matches 5,20,35,50 not 15,30", () => { + expect(matchesCron("5/15 * * * *", at(5))).toBe(true); + expect(matchesCron("5/15 * * * *", at(20))).toBe(true); + expect(matchesCron("5/15 * * * *", at(50))).toBe(true); + expect(matchesCron("5/15 * * * *", at(15))).toBe(false); + expect(matchesCron("5/15 * * * *", at(0))).toBe(false); + }); + + test("plain number matches only itself", () => { + expect(matchesCron("30 * * * *", at(30))).toBe(true); + expect(matchesCron("30 * * * *", at(29))).toBe(false); + }); + + test("plain range without step", () => { + expect(matchesCron("10-20 * * * *", at(15))).toBe(true); + expect(matchesCron("10-20 * * * *", at(21))).toBe(false); + expect(matchesCron("10-20 * * * *", at(9))).toBe(false); + }); + + test("range with step respects both bounds and offset", () => { + expect(matchesCron("0-30/7 * * * *", at(0))).toBe(true); + expect(matchesCron("0-30/7 * * * *", at(28))).toBe(true); + expect(matchesCron("0-30/7 * * * *", at(35))).toBe(false); + expect(matchesCron("5-25/10 * * * *", at(5))).toBe(true); + expect(matchesCron("5-25/10 * * * *", at(15))).toBe(true); + expect(matchesCron("5-25/10 * * * *", at(10))).toBe(false); + }); + + test("invalid step never matches", () => { + expect(matchesCron("*/0 * * * *", at(0))).toBe(false); + expect(matchesCron("*/x * * * *", at(0))).toBe(false); + }); + + test("hour expression: 0 */6 matches 0,6,12,18", () => { + expect(matchesCron("0 */6 * * *", at(0, 0))).toBe(true); + expect(matchesCron("0 */6 * * *", at(0, 6))).toBe(true); + expect(matchesCron("0 */6 * * *", at(0, 18))).toBe(true); + expect(matchesCron("0 */6 * * *", at(0, 3))).toBe(false); + }); + + test("comma list", () => { + expect(matchesCron("1,15 * * * *", at(1))).toBe(true); + expect(matchesCron("1,15 * * * *", at(15))).toBe(true); + expect(matchesCron("1,15 * * * *", at(2))).toBe(false); + }); + + test("wrong field count is invalid", () => { + expect(matchesCron("* * *", at(0))).toBe(false); + expect(matchesCron("* * * * * *", at(0))).toBe(false); + }); +}); diff --git a/apps/api/src/backup/scheduler.ts b/apps/api/src/backup/scheduler.ts index 4b7f340..8ea77e0 100644 --- a/apps/api/src/backup/scheduler.ts +++ b/apps/api/src/backup/scheduler.ts @@ -1,8 +1,8 @@ import { listAllDatabases } from "../db/repo/databases"; import { getBackupStorageSettings } from "../db/repo/settings"; import { BackupOrchestrator } from "./orchestrator"; -import { S3_BACKUP_PREFIX } from "./types"; import type { BackupTarget, StorageConfig } from "./types"; +import { S3_BACKUP_PREFIX } from "./types"; const lastFiredMinute = new Map(); @@ -71,7 +71,7 @@ function toStorageConfig(settings: { return { type: "local", path: settings.path || "/data/backups" }; } -function matchesCron(cron: string, date: Date): boolean { +export function matchesCron(cron: string, date: Date): boolean { const parts = cron.trim().split(/\s+/); if (parts.length !== 5) return false; @@ -90,21 +90,28 @@ function matchField(expr: string, value: number): boolean { if (expr === "*") return true; for (const part of expr.split(",")) { - if (part.includes("-")) { - const [start, end] = part.split("-").map(Number); - if (value >= start && value <= end) return true; - } else if (part.includes("/")) { - const [range, step] = part.split("/"); - const stepNum = parseInt(step, 10); - if (range === "*") { - if (value % stepNum === 0) return true; - } else { - const start = parseInt(range, 10); - if (value >= start && value % stepNum === 0) return true; + if (!part) continue; + const [range, stepStr] = part.split("/"); + const step = stepStr === undefined ? 1 : Number.parseInt(stepStr, 10); + if (!Number.isInteger(step) || step <= 0) continue; + + let start = 0; + let end = Number.POSITIVE_INFINITY; + if (range.includes("-")) { + const [s, e] = range.split("-").map((n) => Number.parseInt(n, 10)); + if (!Number.isInteger(s) || !Number.isInteger(e)) continue; + start = s; + end = e; + } else if (range !== "*") { + const s = Number.parseInt(range, 10); + if (!Number.isInteger(s)) continue; + if (stepStr === undefined) { + if (value === s) return true; + continue; } - } else { - if (parseInt(part, 10) === value) return true; + start = s; } + if (value >= start && value <= end && (value - start) % step === 0) return true; } return false; diff --git a/apps/api/src/db/__tests__/deployment-events.test.ts b/apps/api/src/db/__tests__/deployment-events.test.ts new file mode 100644 index 0000000..49e88ff --- /dev/null +++ b/apps/api/src/db/__tests__/deployment-events.test.ts @@ -0,0 +1,201 @@ +import { afterAll, afterEach, beforeAll, describe, expect, it } from "bun:test"; +import { drizzle } from "drizzle-orm/node-postgres"; +import type { Pool } from "pg"; +import { setDbProvider } from "../db-provider"; +import * as schema from "../schema"; +import { createTestPool, truncateAllTables } from "../test-helper"; + +let pool: Pool; + +const seed = async () => { + await pool.query( + `INSERT INTO projects (id, name, source_type, created_at, updated_at) + VALUES ('proj-ev', 'Event Project', 'git', NOW(), NOW()) ON CONFLICT DO NOTHING`, + ); + await pool.query( + `INSERT INTO deployments (id, project_id, source_type, source_ref, status, branch, commit_sha, created_at, updated_at) + VALUES ('dep-ev-1', 'proj-ev', 'git', 'https://github.com/test/repo.git', 'pending', 'main', 'abc1234567890abcdef', NOW(), NOW()), + ('dep-ev-2', 'proj-ev', 'git', 'https://github.com/test/repo.git', 'pending', 'main', NULL, NOW(), NOW()) + ON CONFLICT DO NOTHING`, + ); +}; + +beforeAll(async () => { + pool = createTestPool(); + const db = drizzle(pool, { schema }); + setDbProvider(async () => db); + await truncateAllTables(pool); + await seed(); +}); + +afterEach(async () => { + await truncateAllTables(pool); + await seed(); +}); + +afterAll(async () => { + try { + await truncateAllTables(pool); + } finally { + await pool.end(); + } +}); + +const eventsFor = async (deploymentId: string) => + ( + await pool.query(`SELECT id, type, message, metadata, sent_at FROM deployment_events WHERE deployment_id = $1`, [ + deploymentId, + ]) + ).rows; + +const deploymentRow = async (deploymentId: string) => + (await pool.query(`SELECT status, failure_reason FROM deployments WHERE id = $1`, [deploymentId])).rows[0]; + +describe("deployment event backbone", () => { + it("records a failure exactly once under 5 concurrent calls", async () => { + const { recordDeploymentFailure } = await import("../repo/deployment-events"); + const results = await Promise.all( + Array.from({ length: 5 }, () => + recordDeploymentFailure({ deploymentId: "dep-ev-1", reason: "build exploded", source: "pipeline" }), + ), + ); + + expect(results.filter((r) => r.claimed)).toHaveLength(1); + expect(results.filter((r) => r.eventId !== null)).toHaveLength(1); + + const events = await eventsFor("dep-ev-1"); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("failed"); + expect(events[0].message).toBe("build exploded"); + + const dep = await deploymentRow("dep-ev-1"); + expect(dep.status).toBe("failed"); + expect(dep.failure_reason).toBe("build exploded"); + }); + + it("keeps the first terminal event across failed -> pending -> failed", async () => { + const { recordDeploymentFailure } = await import("../repo/deployment-events"); + const { updateDeploymentStatus } = await import("../repo/deployments"); + + const first = await recordDeploymentFailure({ + deploymentId: "dep-ev-1", + reason: "first failure", + source: "pipeline", + }); + expect(first.claimed).toBe(true); + + await updateDeploymentStatus("dep-ev-1", "pending"); + + const second = await recordDeploymentFailure({ + deploymentId: "dep-ev-1", + reason: "second failure", + source: "ssh", + }); + expect(second.claimed).toBe(false); + + const events = await eventsFor("dep-ev-1"); + expect(events).toHaveLength(1); + expect(events[0].message).toBe("first failure"); + }); + + it("suppresses failure emails when the deployment was cancelled", async () => { + const { recordDeploymentCancellation, recordDeploymentFailure, listPendingFailureNotificationIds } = await import( + "../repo/deployment-events" + ); + + const cancel = await recordDeploymentCancellation({ + deploymentId: "dep-ev-1", + reason: "Cancelled", + source: "pipeline", + }); + expect(cancel.claimed).toBe(true); + + const fail = await recordDeploymentFailure({ + deploymentId: "dep-ev-1", + reason: "late agent error", + source: "job-channel", + }); + expect(fail.claimed).toBe(false); + + const events = await eventsFor("dep-ev-1"); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("cancelled"); + + const dep = await deploymentRow("dep-ev-1"); + expect(dep.status).toBe("failed"); + expect(dep.failure_reason).toBe("Cancelled"); + + expect(await listPendingFailureNotificationIds()).toHaveLength(0); + }); + + it("does not let the trigger add a failed event after a cancellation", async () => { + const { recordDeploymentCancellation } = await import("../repo/deployment-events"); + await recordDeploymentCancellation({ deploymentId: "dep-ev-1", reason: "Cancelled", source: "pipeline" }); + + await pool.query(`UPDATE deployments SET status = 'pending', failure_reason = NULL WHERE id = 'dep-ev-1'`); + await pool.query(`UPDATE deployments SET status = 'failed', failure_reason = 'Cancelled' WHERE id = 'dep-ev-1'`); + + const events = await eventsFor("dep-ev-1"); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("cancelled"); + }); + + it("trigger inserts a fallback failed event for raw status writes", async () => { + await pool.query( + `UPDATE deployments SET status = 'failed', failure_reason = 'raw write boom' WHERE id = 'dep-ev-2'`, + ); + + const events = await eventsFor("dep-ev-2"); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("failed"); + expect(events[0].message).toBe("raw write boom"); + expect((events[0].metadata as any)?.source).toBe("status-trigger"); + expect(events[0].sent_at).toBeNull(); + }); + + it("claims, marks sent, and stops retrying", async () => { + const { + recordDeploymentFailure, + claimFailureNotification, + markFailureNotificationSent, + listPendingFailureNotificationIds, + } = await import("../repo/deployment-events"); + + const { claimed, eventId } = await recordDeploymentFailure({ + deploymentId: "dep-ev-1", + reason: "mail me", + source: "pipeline", + }); + expect(claimed).toBe(true); + + const pending = await listPendingFailureNotificationIds(); + expect(pending).toContain(eventId); + + const ctx = await claimFailureNotification(eventId!); + expect(ctx).not.toBeNull(); + expect(ctx!.deploymentId).toBe("dep-ev-1"); + expect(ctx!.projectName).toBe("Event Project"); + expect(ctx!.failureReason).toBe("mail me"); + expect(ctx!.commitSha).toBe("abc1234567890abcdef"); + expect(ctx!.attempt).toBe(1); + + await markFailureNotificationSent(eventId!); + expect(await claimFailureNotification(eventId!)).toBeNull(); + expect(await listPendingFailureNotificationIds()).not.toContain(eventId); + }); + + it("caps delivery attempts at 3", async () => { + const { recordDeploymentFailure, claimFailureNotification } = await import("../repo/deployment-events"); + + const { eventId } = await recordDeploymentFailure({ + deploymentId: "dep-ev-2", + reason: "flaky smtp", + source: "pipeline", + }); + + expect(await claimFailureNotification(eventId!)).not.toBeNull(); + expect(await claimFailureNotification(eventId!)).not.toBeNull(); + expect(await claimFailureNotification(eventId!)).not.toBeNull(); + expect(await claimFailureNotification(eventId!)).toBeNull(); + }); +}); diff --git a/apps/api/src/db/__tests__/shared-env-links-runner.ts b/apps/api/src/db/__tests__/shared-env-links-runner.ts index 61dbcea..8b8fe62 100644 --- a/apps/api/src/db/__tests__/shared-env-links-runner.ts +++ b/apps/api/src/db/__tests__/shared-env-links-runner.ts @@ -1,14 +1,14 @@ -import { drizzle } from "drizzle-orm/node-postgres"; -import { Pool } from "pg"; import { randomUUID } from "node:crypto"; import { eq } from "drizzle-orm"; +import { drizzle } from "drizzle-orm/node-postgres"; +import { Pool } from "pg"; import { setDbProvider } from "../db-provider"; -import * as schema from "../schema"; import { linkSharedEnvVarsToProject, listLinkedSharedEnvVars, unlinkSharedEnvVarFromProject, } from "../repo/shared-env-vars"; +import * as schema from "../schema"; const TEST_DATABASE_URL = process.env.TEST_DATABASE_URL ?? "postgresql://dequel:dequel@localhost:5433/dequel"; const pool = new Pool({ connectionString: TEST_DATABASE_URL }); @@ -45,7 +45,9 @@ try { "created_at" timestamp DEFAULT now() ) `); - await pool.query(`CREATE UNIQUE INDEX IF NOT EXISTS "project_shared_env_links_project_shared_idx" ON "project_shared_env_links" ("project_id", "shared_env_var_id")`); + await pool.query( + `CREATE UNIQUE INDEX IF NOT EXISTS "project_shared_env_links_project_shared_idx" ON "project_shared_env_links" ("project_id", "shared_env_var_id")`, + ); await cleanup(); const projectId = `test-sel-${randomUUID().slice(0, 8)}`; diff --git a/apps/api/src/db/migrate.ts b/apps/api/src/db/migrate.ts index 5843151..8707c6d 100644 --- a/apps/api/src/db/migrate.ts +++ b/apps/api/src/db/migrate.ts @@ -2,7 +2,6 @@ import { migrate as drizzleMigrate } from "drizzle-orm/node-postgres/migrator"; import { config } from "../utils/config"; import { getDb } from "./client"; import { getGithubIntegration, setGithubIntegration } from "./repo/github"; -import { getSmtpSettings, upsertSmtpSettings } from "./repo/settings"; export const migrate = async () => { const db = await getDb(); @@ -51,17 +50,4 @@ const seedFromConfig = async () => { console.log("[Config] Seeded GitHub integration from config file"); } } - if (config.smtpHost) { - const existing = await getSmtpSettings(); - if (!existing) { - await upsertSmtpSettings({ - host: config.smtpHost, - port: config.smtpPort, - user: config.smtpUser, - pass: config.smtpPass, - fromAddress: config.smtpFrom, - }); - console.log("[Config] Seeded SMTP settings from config file"); - } - } }; diff --git a/apps/api/src/db/migrations/0034_deployment_event_backbone.sql b/apps/api/src/db/migrations/0034_deployment_event_backbone.sql new file mode 100644 index 0000000..019c40f --- /dev/null +++ b/apps/api/src/db/migrations/0034_deployment_event_backbone.sql @@ -0,0 +1,29 @@ +ALTER TABLE deployment_events ADD COLUMN sent_at timestamptz; +ALTER TABLE deployment_events ADD COLUMN attempts integer NOT NULL DEFAULT 0; + +UPDATE deployment_events SET sent_at = now() WHERE type = 'failed'; + +DELETE FROM deployment_events a USING deployment_events b + WHERE a.deployment_id = b.deployment_id AND a.type = b.type AND a.ctid < b.ctid; + +CREATE UNIQUE INDEX udep_events_failed ON deployment_events (deployment_id) WHERE type = 'failed'; +CREATE UNIQUE INDEX udep_events_cancelled ON deployment_events (deployment_id) WHERE type = 'cancelled'; + +CREATE FUNCTION deployment_terminal_event_fallback() RETURNS trigger AS $$ +BEGIN + IF NOT EXISTS (SELECT 1 FROM deployment_events + WHERE deployment_id = NEW.id AND type IN ('failed','cancelled')) THEN + INSERT INTO deployment_events (id, deployment_id, type, message, metadata) + VALUES (gen_random_uuid()::text, NEW.id, 'failed', NEW.failure_reason, + jsonb_build_object('source', 'status-trigger')) + ON CONFLICT DO NOTHING; + END IF; + RETURN NEW; +END $$ LANGUAGE plpgsql; + +CREATE TRIGGER trg_deployment_terminal_event +AFTER UPDATE OF status ON deployments +FOR EACH ROW WHEN (NEW.status = 'failed' AND OLD.status IS DISTINCT FROM 'failed') +EXECUTE FUNCTION deployment_terminal_event_fallback(); + +DELETE FROM alerts WHERE type = 'error_rate'; diff --git a/apps/api/src/db/migrations/0035_deployment_fk_cascade.sql b/apps/api/src/db/migrations/0035_deployment_fk_cascade.sql new file mode 100644 index 0000000..4ddc59e --- /dev/null +++ b/apps/api/src/db/migrations/0035_deployment_fk_cascade.sql @@ -0,0 +1,11 @@ +ALTER TABLE "agent_jobs" DROP CONSTRAINT IF EXISTS "agent_jobs_deployment_id_deployments_id_fk"; +ALTER TABLE "agent_jobs" ADD CONSTRAINT "agent_jobs_deployment_id_deployments_id_fk" + FOREIGN KEY ("deployment_id") REFERENCES "deployments"("id") ON DELETE CASCADE; + +ALTER TABLE "deployment_logs" DROP CONSTRAINT IF EXISTS "deployment_logs_deployment_id_deployments_id_fk"; +ALTER TABLE "deployment_logs" ADD CONSTRAINT "deployment_logs_deployment_id_deployments_id_fk" + FOREIGN KEY ("deployment_id") REFERENCES "deployments"("id") ON DELETE CASCADE; + +ALTER TABLE "deployment_events" DROP CONSTRAINT IF EXISTS "deployment_events_deployment_id_deployments_id_fk"; +ALTER TABLE "deployment_events" ADD CONSTRAINT "deployment_events_deployment_id_deployments_id_fk" + FOREIGN KEY ("deployment_id") REFERENCES "deployments"("id") ON DELETE CASCADE; diff --git a/apps/api/src/db/migrations/meta/_journal.json b/apps/api/src/db/migrations/meta/_journal.json index 6497c94..c8aa637 100644 --- a/apps/api/src/db/migrations/meta/_journal.json +++ b/apps/api/src/db/migrations/meta/_journal.json @@ -71,6 +71,20 @@ "when": 1788600000000, "tag": "0033_add_system_backup_settings", "breakpoints": true + }, + { + "idx": 10, + "version": "7", + "when": 1788600001000, + "tag": "0034_deployment_event_backbone", + "breakpoints": true + }, + { + "idx": 11, + "version": "7", + "when": 1790424535143, + "tag": "0035_deployment_fk_cascade", + "breakpoints": true } ] -} +} \ No newline at end of file diff --git a/apps/api/src/db/repo/deployment-events.ts b/apps/api/src/db/repo/deployment-events.ts index 528dbad..7ffd3a6 100644 --- a/apps/api/src/db/repo/deployment-events.ts +++ b/apps/api/src/db/repo/deployment-events.ts @@ -1,7 +1,12 @@ import { randomUUID } from "node:crypto"; -import { eq } from "drizzle-orm"; +import { and, asc, desc, eq, inArray, isNull, lt, sql } from "drizzle-orm"; +import { emitDeploymentFailed } from "../../events"; +import type { FailureNotificationContext, RecordFailureInput, RecordFailureOutcome } from "../../types"; import { getDb } from "../db-provider"; -import { deploymentEvents } from "../schema"; +import { deploymentEvents, deployments, projects } from "../schema"; +import { applyStatusUpdate } from "./deployments"; + +const TERMINAL_TYPES = ["failed", "cancelled"]; export const createDeploymentEvent = async (input: { deploymentId: string; @@ -33,3 +38,180 @@ export const listDeploymentEvents = async (deploymentId: string) => { .orderBy(deploymentEvents.createdAt) .execute(); }; + +const recordTerminal = async ( + input: RecordFailureInput, + type: "failed" | "cancelled", +): Promise => { + const db = await getDb(); + const eventId = await db.transaction(async (tx) => { + const [dep] = await tx + .select({ id: deployments.id }) + .from(deployments) + .where(eq(deployments.id, input.deploymentId)) + .for("update") + .execute(); + if (!dep) return null; + + const [existing] = await tx + .select({ id: deploymentEvents.id }) + .from(deploymentEvents) + .where(and(eq(deploymentEvents.deploymentId, input.deploymentId), inArray(deploymentEvents.type, TERMINAL_TYPES))) + .limit(1) + .execute(); + if (existing) return null; + + const [inserted] = await tx + .insert(deploymentEvents) + .values({ + id: randomUUID(), + deploymentId: input.deploymentId, + type, + message: input.reason, + metadata: { source: input.source }, + }) + .onConflictDoNothing() + .returning({ id: deploymentEvents.id }) + .execute(); + if (!inserted) return null; + + await applyStatusUpdate(tx, input.deploymentId, "failed", { failureReason: input.reason }); + return inserted.id; + }); + + if (eventId && type === "failed") { + emitDeploymentFailed({ eventId, deploymentId: input.deploymentId }); + } + return { claimed: eventId !== null, eventId }; +}; + +export const recordDeploymentFailure = (input: RecordFailureInput): Promise => + recordTerminal({ ...input, cancel: false }, "failed"); + +export const recordDeploymentCancellation = (input: RecordFailureInput): Promise => + recordTerminal({ ...input, cancel: true }, "cancelled"); + +export const claimFailureNotification = async (eventId: string): Promise => { + const db = await getDb(); + const [claimed] = await db + .update(deploymentEvents) + .set({ attempts: sql`${deploymentEvents.attempts} + 1` }) + .where( + and( + eq(deploymentEvents.id, eventId), + eq(deploymentEvents.type, "failed"), + isNull(deploymentEvents.sentAt), + lt(deploymentEvents.attempts, 3), + ), + ) + .returning({ id: deploymentEvents.id }) + .execute(); + if (!claimed) return null; + + const [row] = await db + .select({ + eventId: deploymentEvents.id, + deploymentId: deploymentEvents.deploymentId, + attempts: deploymentEvents.attempts, + message: deploymentEvents.message, + projectId: deployments.projectId, + sourceRef: deployments.sourceRef, + commitSha: deployments.commitSha, + finishedAt: deployments.finishedAt, + projectName: projects.name, + }) + .from(deploymentEvents) + .innerJoin(deployments, eq(deployments.id, deploymentEvents.deploymentId)) + .leftJoin(projects, eq(projects.id, deployments.projectId)) + .where(eq(deploymentEvents.id, eventId)) + .execute(); + if (!row) return null; + + return { + eventId: row.eventId, + deploymentId: row.deploymentId, + projectId: row.projectId, + projectName: row.projectName ?? row.sourceRef, + failureReason: row.message, + commitSha: row.commitSha, + sourceRef: row.sourceRef, + finishedAt: row.finishedAt ? new Date(row.finishedAt).toISOString() : null, + attempt: row.attempts, + }; +}; + +export const markFailureNotificationSent = async (eventId: string): Promise => { + const db = await getDb(); + await db.update(deploymentEvents).set({ sentAt: new Date() }).where(eq(deploymentEvents.id, eventId)).execute(); +}; + +export const listPendingFailureNotificationIds = async (limit = 20): Promise => { + const db = await getDb(); + const rows = await db + .select({ id: deploymentEvents.id }) + .from(deploymentEvents) + .where(and(eq(deploymentEvents.type, "failed"), isNull(deploymentEvents.sentAt), lt(deploymentEvents.attempts, 3))) + .orderBy(asc(deploymentEvents.createdAt)) + .limit(limit) + .execute(); + return rows.map((r) => r.id); +}; + +export interface ProjectEventRow { + id: string; + deploymentId: string; + type: string; + message: string | null; + source: string | null; + createdAt: string; + deploymentStatus: string; + commitSha: string | null; + sourceRef: string; + finishedAt: string | null; + sentAt: string | null; +} + +export const listProjectEvents = async ( + projectId: string, + sinceIso: string, + limit = 200, +): Promise => { + const db = await getDb(); + const rows = await db + .select({ + id: deploymentEvents.id, + deploymentId: deploymentEvents.deploymentId, + type: deploymentEvents.type, + message: deploymentEvents.message, + metadata: deploymentEvents.metadata, + createdAt: deploymentEvents.createdAt, + deploymentStatus: deployments.status, + commitSha: deployments.commitSha, + sourceRef: deployments.sourceRef, + finishedAt: deployments.finishedAt, + sentAt: deploymentEvents.sentAt, + }) + .from(deploymentEvents) + .innerJoin(deployments, eq(deployments.id, deploymentEvents.deploymentId)) + .where(and(eq(deployments.projectId, projectId), sql`${deploymentEvents.createdAt} >= ${sinceIso}`)) + .orderBy(desc(deploymentEvents.createdAt)) + .limit(limit) + .execute(); + + return rows.map((r) => ({ + id: r.id, + deploymentId: r.deploymentId, + type: r.type, + message: r.message, + source: + r.metadata && typeof r.metadata === "object" && r.metadata !== null + ? ((r.metadata as Record).source as string | null) + : null, + createdAt: new Date(r.createdAt).toISOString(), + deploymentStatus: r.deploymentStatus, + commitSha: r.commitSha, + sourceRef: r.sourceRef, + finishedAt: r.finishedAt ? new Date(r.finishedAt).toISOString() : null, + sentAt: r.sentAt ? new Date(r.sentAt).toISOString() : null, + })); +}; diff --git a/apps/api/src/db/repo/deployments.ts b/apps/api/src/db/repo/deployments.ts index 4b5a135..d95ef5f 100644 --- a/apps/api/src/db/repo/deployments.ts +++ b/apps/api/src/db/repo/deployments.ts @@ -91,13 +91,15 @@ const ACTIVE_STATUSES: DeploymentStatus[] = ["pending", "building", "deploying"] const STAMP_FINISHED_UNCONDITIONALLY: DeploymentStatus[] = ["running", "failed"]; -export const updateDeploymentStatus = async ( +type Tx = Parameters>["transaction"]>[0]>[0]; + +export const applyStatusUpdate = async ( + tx: Tx | Awaited>, id: string, status: DeploymentStatus, patch: Partial> = {}, ) => { - const db = await getDb(); - const [existing] = await db + const [existing] = await tx .select({ finishedAt: deployments.finishedAt }) .from(deployments) .where(eq(deployments.id, id)) @@ -115,7 +117,16 @@ export const updateDeploymentStatus = async ( updates.finishedAt = now(); } } - await db.update(deployments).set(updates).where(eq(deployments.id, id)).execute(); + await tx.update(deployments).set(updates).where(eq(deployments.id, id)).execute(); +}; + +export const updateDeploymentStatus = async ( + id: string, + status: DeploymentStatus, + patch: Partial> = {}, +) => { + const db = await getDb(); + await applyStatusUpdate(db, id, status, patch); }; export const deleteDeploymentAndLogs = async (id: string): Promise => { diff --git a/apps/api/src/db/repo/index.ts b/apps/api/src/db/repo/index.ts index fa4369b..cd908bc 100644 --- a/apps/api/src/db/repo/index.ts +++ b/apps/api/src/db/repo/index.ts @@ -29,7 +29,16 @@ export { updateDatabaseSettings, updateDatabaseStatus, } from "./databases"; -export { createDeploymentEvent, listDeploymentEvents } from "./deployment-events"; +export { + claimFailureNotification, + createDeploymentEvent, + listDeploymentEvents, + listPendingFailureNotificationIds, + listProjectEvents, + markFailureNotificationSent, + recordDeploymentCancellation, + recordDeploymentFailure, +} from "./deployment-events"; export { appendLog, countDeployments, diff --git a/apps/api/src/db/repo/shared-env-vars.ts b/apps/api/src/db/repo/shared-env-vars.ts index 44b6c0d..ee34d54 100644 --- a/apps/api/src/db/repo/shared-env-vars.ts +++ b/apps/api/src/db/repo/shared-env-vars.ts @@ -136,12 +136,8 @@ export const unlinkSharedEnvVarFromProject = async (projectId: string, sharedEnv .execute(); if (links.length === 0) return false; return ( - getRowsAffected( - await db - .delete(projectSharedEnvLinks) - .where(eq(projectSharedEnvLinks.id, links[0].id)) - .execute(), - ) > 0 + getRowsAffected(await db.delete(projectSharedEnvLinks).where(eq(projectSharedEnvLinks.id, links[0].id)).execute()) > + 0 ); }; diff --git a/apps/api/src/db/schema.ts b/apps/api/src/db/schema.ts index 93a085a..eccf8b7 100644 --- a/apps/api/src/db/schema.ts +++ b/apps/api/src/db/schema.ts @@ -104,6 +104,8 @@ export const deploymentEvents = pgTable( type: text().notNull(), message: text(), metadata: jsonb("metadata"), + sentAt: timestamp("sent_at", { withTimezone: true }), + attempts: integer().notNull().default(0), createdAt: timestamp("created_at", { withTimezone: true }).notNull().defaultNow(), }, (table) => [ diff --git a/apps/api/src/events.ts b/apps/api/src/events.ts new file mode 100644 index 0000000..4530ac7 --- /dev/null +++ b/apps/api/src/events.ts @@ -0,0 +1,23 @@ +export interface DeploymentFailedSignal { + eventId: string; + deploymentId: string; +} + +type Handler = (signal: DeploymentFailedSignal) => void; + +const handlers = new Set(); + +export const onDeploymentFailed = (handler: Handler): (() => void) => { + handlers.add(handler); + return () => handlers.delete(handler); +}; + +export const emitDeploymentFailed = (signal: DeploymentFailedSignal): void => { + for (const handler of handlers) { + try { + handler(signal); + } catch (err) { + console.error("[Events] deployment-failed handler error:", err); + } + } +}; diff --git a/apps/api/src/executors/agent.ts b/apps/api/src/executors/agent.ts index 8a9bb09..9c69cbf 100644 --- a/apps/api/src/executors/agent.ts +++ b/apps/api/src/executors/agent.ts @@ -1,4 +1,5 @@ import type { Deployment, Project, Server } from "../types"; +import { CANCELLED_FAILURE_REASON } from "../utils/failure-outcome"; import { routeNamesFor } from "../utils/routes"; import type { DeploymentExecutor, @@ -106,10 +107,14 @@ export const agentExecutor: DeploymentExecutor = { }, async cancel({ deployment }: ExecutorCancelInput) { - const { cancelAgentJobsByDeploymentId, updateDeploymentStatus, appendLog } = await getRepo(); + const { cancelAgentJobsByDeploymentId, recordDeploymentCancellation, appendLog } = await getRepo(); if (deployment.status !== "pending" && deployment.status !== "building") return; await cancelAgentJobsByDeploymentId(deployment.id); - await updateDeploymentStatus(deployment.id, "failed", { failureReason: "Cancelled" }); + await recordDeploymentCancellation({ + deploymentId: deployment.id, + reason: CANCELLED_FAILURE_REASON, + source: "agent", + }); await appendLog(deployment.id, "system", "Deployment cancelled by user"); }, }; diff --git a/apps/api/src/executors/ssh.ts b/apps/api/src/executors/ssh.ts index dcac72b..59f4881 100644 --- a/apps/api/src/executors/ssh.ts +++ b/apps/api/src/executors/ssh.ts @@ -1,6 +1,7 @@ import { summarizeDeploymentError } from "../orchestrator/deployment-errors"; import type { Deployment, Project, Server } from "../types"; import { config } from "../utils/config"; +import { CANCELLED_FAILURE_REASON } from "../utils/failure-outcome"; import { removeRemoteCaddyRoute, runRemoteScript, syncRemoteCaddyRoute } from "../utils/ssh"; import { emitLog } from "./logging"; import { buildRemoteDeployScript, parseRemoteBuildResult } from "./ssh-build-script"; @@ -254,10 +255,10 @@ const deployComposeRemote = async (deployment: Deployment, project: Project, ser }; const markFailed = async (deploymentId: string, error: unknown) => { - const { updateDeploymentStatus } = await getRepo(); + const { recordDeploymentFailure } = await getRepo(); const message = summarizeDeploymentError(error); await emitLog(deploymentId, "system", `Deployment failed: ${message}`); - await updateDeploymentStatus(deploymentId, "failed", { failureReason: message }); + await recordDeploymentFailure({ deploymentId, reason: message, source: "ssh" }); }; export const sshExecutor: DeploymentExecutor = { @@ -268,7 +269,11 @@ export const sshExecutor: DeploymentExecutor = { if (!project) throw new Error("Deployment requires a project"); if (project.buildType === "compose") { - await deployComposeRemote(deployment, project, server); + try { + await deployComposeRemote(deployment, project, server); + } catch (error) { + await markFailed(deployment.id, error); + } return; } @@ -353,7 +358,8 @@ export const sshExecutor: DeploymentExecutor = { } catch (error) { const message = summarizeDeploymentError(error); await emitLog(deployment.id, "system", `Rollback failed: ${message}`); - await updateDeploymentStatus(deployment.id, "failed", { failureReason: message }); + const { recordDeploymentFailure } = await getRepo(); + await recordDeploymentFailure({ deploymentId: deployment.id, reason: message, source: "rollback" }); throw error; } }, @@ -400,9 +406,13 @@ export const sshExecutor: DeploymentExecutor = { }, async cancel({ deployment }: ExecutorCancelInput) { - const { updateDeploymentStatus } = await getRepo(); + const { recordDeploymentCancellation } = await getRepo(); if (deployment.status !== "pending" && deployment.status !== "building") return; - await updateDeploymentStatus(deployment.id, "failed", { failureReason: "Cancelled" }); + await recordDeploymentCancellation({ + deploymentId: deployment.id, + reason: CANCELLED_FAILURE_REASON, + source: "ssh", + }); await emitLog(deployment.id, "system", "Deployment cancelled by user (remote build may continue on the server)"); }, }; diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index cf5dc6a..8297ea1 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -10,6 +10,7 @@ import { migrate } from "./db/migrate"; import { ensureLocalServer } from "./db/repo"; import { deployments } from "./db/schema"; import { alertEvaluator } from "./monitoring/evaluator"; +import { startFailureNotifier } from "./monitoring/failure-notifier"; import { orchestrator } from "./orchestrator"; import { startBuildCleanup } from "./orchestrator/cleanup"; import { startFailoverMonitor } from "./orchestrator/failover"; @@ -36,6 +37,7 @@ const bootstrap = async () => { serverManager.start(); startDomainPolling(); alertEvaluator.start(); + startFailureNotifier(); startBuildCleanup(); startFailoverMonitor(); startDatabaseMonitoring(); diff --git a/apps/api/src/monitoring/__tests__/alert-guard.test.ts b/apps/api/src/monitoring/__tests__/alert-guard.test.ts new file mode 100644 index 0000000..5b4fcd4 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/alert-guard.test.ts @@ -0,0 +1,70 @@ +import { describe, expect, it } from "bun:test"; +import { scalingGuard } from "../alert-guard"; + +const enabled = (maxReplicas = 5) => ({ enabled: true, maxReplicas }); + +describe("scalingGuard", () => { + it("does not touch downtime or unknown alert types", () => { + expect(scalingGuard("downtime", { policy: null, cpuLimit: null, currentReplicas: null })).toEqual({ + suppress: false, + suggestion: null, + }); + }); + + it("fires with an enable prompt when no policy exists", () => { + expect(scalingGuard("cpu", { policy: null, cpuLimit: 1, currentReplicas: 1 })).toEqual({ + suppress: false, + suggestion: { kind: "enable_autoscaling" }, + }); + }); + + it("fires with an enable prompt when the policy is disabled", () => { + expect( + scalingGuard("memory", { + policy: { enabled: false, maxReplicas: 5 }, + cpuLimit: 1, + currentReplicas: 1, + }), + ).toEqual({ suppress: false, suggestion: { kind: "enable_autoscaling" } }); + }); + + it("fires with an enable prompt when no cpu limit is set", () => { + expect(scalingGuard("cpu", { policy: enabled(), cpuLimit: null, currentReplicas: 1 })).toEqual({ + suppress: false, + suggestion: { kind: "enable_autoscaling" }, + }); + expect(scalingGuard("cpu", { policy: enabled(), cpuLimit: 0, currentReplicas: 1 })).toEqual({ + suppress: false, + suggestion: { kind: "enable_autoscaling" }, + }); + }); + + it("fires with an increase-replicas prompt when at the ceiling", () => { + expect(scalingGuard("cpu", { policy: enabled(3), cpuLimit: 1, currentReplicas: 3 })).toEqual({ + suppress: false, + suggestion: { kind: "increase_max_replicas", current: 3, maxReplicas: 3 }, + }); + }); + + it("treats an unknown replica count as a single replica", () => { + expect(scalingGuard("cpu", { policy: enabled(1), cpuLimit: 1, currentReplicas: null })).toEqual({ + suppress: false, + suggestion: { kind: "increase_max_replicas", current: 1, maxReplicas: 1 }, + }); + expect(scalingGuard("cpu", { policy: enabled(5), cpuLimit: 1, currentReplicas: null })).toEqual({ + suppress: true, + suggestion: null, + }); + }); + + it("suppresses when autoscaling is on and replicas are below the ceiling", () => { + expect(scalingGuard("cpu", { policy: enabled(5), cpuLimit: 1, currentReplicas: 1 })).toEqual({ + suppress: true, + suggestion: null, + }); + expect(scalingGuard("memory", { policy: enabled(5), cpuLimit: 1, currentReplicas: 4 })).toEqual({ + suppress: true, + suggestion: null, + }); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/health.test.ts b/apps/api/src/monitoring/__tests__/health.test.ts new file mode 100644 index 0000000..3abf9e1 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/health.test.ts @@ -0,0 +1,104 @@ +import { describe, expect, it } from "bun:test"; +import { evaluateHealth } from "../health"; + +const base = { + server: null as { status: string; lastHeartbeatAgeMs: number | null } | null, + runningDeployments: 1, + hasDeployments: true, + routeErrors: [] as string[], + http: { available: true, errorRate: 0 }, +}; + +const check = (result: ReturnType, name: string) => result.checks.find((c) => c.name === name)!; + +describe("evaluateHealth", () => { + it("is healthy when everything is fine locally", () => { + const result = evaluateHealth(base); + expect(result.overall).toBe("healthy"); + expect(result.checks).toHaveLength(4); + expect(check(result, "server").status).toBe("ok"); + expect(check(result, "containers").status).toBe("ok"); + expect(check(result, "ingress").status).toBe("ok"); + expect(check(result, "http").status).toBe("ok"); + }); + + it("treats a missing server row as local ok", () => { + const result = evaluateHealth({ ...base, server: null }); + expect(check(result, "server").status).toBe("ok"); + }); + + it("fails on a stale remote heartbeat", () => { + const result = evaluateHealth({ ...base, server: { status: "connected", lastHeartbeatAgeMs: 120_000 } }); + expect(check(result, "server").status).toBe("fail"); + expect(check(result, "server").detail).toContain("Heartbeat stale"); + expect(result.overall).toBe("down"); + }); + + it("fails when a server never heartbeated", () => { + const result = evaluateHealth({ ...base, server: { status: "pending", lastHeartbeatAgeMs: null } }); + expect(check(result, "server").status).toBe("fail"); + expect(result.overall).toBe("down"); + }); + + it("passes on a fresh heartbeat", () => { + const result = evaluateHealth({ ...base, server: { status: "connected", lastHeartbeatAgeMs: 5_000 } }); + expect(check(result, "server").status).toBe("ok"); + }); + + it("warns when there are no deployments yet", () => { + const result = evaluateHealth({ ...base, hasDeployments: false, runningDeployments: 0 }); + expect(check(result, "containers").status).toBe("warn"); + expect(check(result, "containers").detail).toBe("No deployments yet"); + expect(result.overall).toBe("degraded"); + }); + + it("fails when deployments exist but none are running", () => { + const result = evaluateHealth({ ...base, runningDeployments: 0 }); + expect(check(result, "containers").status).toBe("fail"); + expect(result.overall).toBe("down"); + }); + + it("warns on route errors with the first error as detail", () => { + const result = evaluateHealth({ ...base, routeErrors: ["upstream unreachable", "second"] }); + expect(check(result, "ingress").status).toBe("warn"); + expect(check(result, "ingress").detail).toBe("upstream unreachable"); + expect(result.overall).toBe("degraded"); + }); + + it("is unknown when Loki is unavailable", () => { + const result = evaluateHealth({ ...base, http: { available: false, errorRate: null } }); + expect(check(result, "http").status).toBe("unknown"); + expect(check(result, "http").detail).toBe("Loki unavailable"); + expect(result.overall).toBe("degraded"); + }); + + it("is unknown when there is no traffic", () => { + const result = evaluateHealth({ ...base, http: { available: true, errorRate: null } }); + expect(check(result, "http").status).toBe("unknown"); + expect(check(result, "http").detail).toBe("No traffic"); + expect(result.overall).toBe("degraded"); + }); + + it("warns above a 5% error rate", () => { + const result = evaluateHealth({ ...base, http: { available: true, errorRate: 0.06 } }); + expect(check(result, "http").status).toBe("warn"); + expect(check(result, "http").detail).toBe("6.0% error rate"); + expect(result.overall).toBe("degraded"); + }); + + it("stays ok at or below a 5% error rate", () => { + const result = evaluateHealth({ ...base, http: { available: true, errorRate: 0.05 } }); + expect(check(result, "http").status).toBe("ok"); + expect(result.overall).toBe("healthy"); + }); + + it("lets fail win over warn for overall status", () => { + const result = evaluateHealth({ + ...base, + runningDeployments: 0, + routeErrors: ["broken"], + http: { available: true, errorRate: 0.5 }, + }); + expect(result.overall).toBe("down"); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/http-error-rate.test.ts b/apps/api/src/monitoring/__tests__/http-error-rate.test.ts new file mode 100644 index 0000000..15e5038 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/http-error-rate.test.ts @@ -0,0 +1,65 @@ +import { describe, expect, it } from "bun:test"; +import { parseHttpErrorRate } from "../http-error-rate"; + +const sample = [ + { metric: { status: "200" }, value: ["1758000000", "480"] }, + { metric: { status: "204" }, value: ["1758000000", "20"] }, + { metric: { status: "404" }, value: ["1758000000", "12"] }, + { metric: { status: "500" }, value: ["1758000000", "8"] }, +]; + +describe("parseHttpErrorRate", () => { + it("derives totals and rate from status buckets", () => { + const result = parseHttpErrorRate(sample, 3600); + expect(result.available).toBe(true); + expect(result.windowSeconds).toBe(3600); + expect(result.totalRequests).toBe(520); + expect(result.errorRequests).toBe(20); + expect(result.errorRate).toBeCloseTo(20 / 520, 10); + expect(result.byStatus).toHaveLength(4); + }); + + it("returns errorRate null on no traffic", () => { + const result = parseHttpErrorRate([], 300); + expect(result.available).toBe(true); + expect(result.totalRequests).toBe(0); + expect(result.errorRate).toBeNull(); + }); + + it("marks unavailable payloads instead of reporting 0", () => { + for (const payload of [null, undefined, "nope", { not: "an array" }]) { + const result = parseHttpErrorRate(payload, 300); + expect(result.available).toBe(false); + expect(result.errorRate).toBeNull(); + expect(result.totalRequests).toBe(0); + expect(result.byStatus).toHaveLength(0); + } + }); + + it("filters malformed bucket entries", () => { + const result = parseHttpErrorRate( + [ + { metric: { status: "200" }, value: ["1", "10"] }, + { metric: {}, value: ["1", "5"] }, + { metric: { status: "201" }, value: ["1", "not-a-number"] }, + null, + ], + 60, + ); + expect(result.available).toBe(true); + expect(result.totalRequests).toBe(10); + expect(result.byStatus).toHaveLength(1); + }); + + it("counts 4xx and 5xx as errors but not 3xx", () => { + const result = parseHttpErrorRate( + [ + { metric: { status: "301" }, value: ["1", "100"] }, + { metric: { status: "404" }, value: ["1", "1"] }, + ], + 60, + ); + expect(result.errorRequests).toBe(1); + expect(result.errorRate).toBeCloseTo(1 / 101, 10); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/incident-policy.test.ts b/apps/api/src/monitoring/__tests__/incident-policy.test.ts new file mode 100644 index 0000000..723d857 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/incident-policy.test.ts @@ -0,0 +1,143 @@ +import { describe, expect, it } from "bun:test"; +import { type Incident, DEFAULT_POLICY, backoffMs, commitAfterSend, decide } from "../incident-policy"; + +const P = DEFAULT_POLICY; +const MIN = 60_000; + +const breach = (state: Incident | null, now: number) => decide(state, { kind: "breach" }, now, P); +const clear = (state: Incident | null, now: number) => decide(state, { kind: "clear" }, now, P); +const noData = (state: Incident | null, now: number) => decide(state, { kind: "no_data" }, now, P); + +const firstSend = (now: number): Incident => { + const d = breach(null, now); + if (d.kind !== "send") throw new Error("expected send"); + return d.next; +}; + +describe("backoffMs", () => { + it("quadruples per send and clamps at the cap", () => { + expect(backoffMs(1, P)).toBe(5 * MIN); + expect(backoffMs(2, P)).toBe(20 * MIN); + expect(backoffMs(3, P)).toBe(80 * MIN); + expect(backoffMs(4, P)).toBe(320 * MIN); + expect(backoffMs(5, P)).toBe(P.capMs); + expect(backoffMs(99, P)).toBe(P.capMs); + }); +}); + +describe("decide — breach", () => { + it("first breach sends immediately with the incident opened now", () => { + const d = breach(null, 1_000); + expect(d.kind).toBe("send"); + if (d.kind !== "send") return; + expect(d.next).toEqual({ + status: "active", + since: 1_000, + sends: 1, + nextDueAt: 1_000 + 5 * MIN, + clearSince: null, + }); + }); + + it("does not send again before the backoff is due", () => { + const state = firstSend(0); + expect(breach(state, 4 * MIN + 59_000).kind).toBe("noop"); + }); + + it("sends the next rung once due, quadrupling the gap", () => { + const state = firstSend(0); + const d = breach(state, 5 * MIN); + expect(d.kind).toBe("send"); + if (d.kind !== "send") return; + expect(d.next.sends).toBe(2); + expect(d.next.nextDueAt).toBe(5 * MIN + 20 * MIN); + }); + + it("sends exactly 5 emails over an 8h sustained breach", () => { + let state: Incident | null = null; + let sends = 0; + for (let minute = 0; minute <= 8 * 60; minute++) { + const d = breach(state, minute * MIN); + if (d.kind === "send") { + sends++; + state = d.next; + } else if (d.kind === "commit") { + state = d.next; + } + } + expect(sends).toBe(5); + }); +}); + +describe("decide — no_data and clear", () => { + it("no_data never fires and never mutates state", () => { + const state = firstSend(0); + expect(noData(state, 10 * MIN).kind).toBe("noop"); + expect(noData(null, 10 * MIN).kind).toBe("noop"); + }); + + it("clear with no incident is a noop", () => { + expect(clear(null, 0).kind).toBe("noop"); + }); + + it("clear starts a recovery grace window without sending", () => { + const state = firstSend(0); + const d = clear(state, 2 * MIN); + expect(d.kind).toBe("commit"); + if (d.kind !== "commit" || !d.next) return; + expect(d.next.clearSince).toBe(2 * MIN); + expect(d.next.sends).toBe(1); + }); + + it("clear inside the grace keeps the incident, past the grace deletes it", () => { + const state = { ...firstSend(0), clearSince: 1 * MIN } as Incident; + const inside = clear(state, 1 * MIN + P.recoveryGraceMs - 1_000); + expect(inside.kind).toBe("noop"); + const past = clear(state, 1 * MIN + P.recoveryGraceMs); + expect(past.kind).toBe("commit"); + if (past.kind !== "commit") return; + expect(past.next).toBeNull(); + }); + + it("a re-breach inside the grace cancels recovery and keeps the schedule", () => { + const state = { ...firstSend(0), clearSince: 2 * MIN } as Incident; + const d = breach(state, 3 * MIN); + expect(d.kind).toBe("commit"); + if (d.kind !== "commit" || !d.next) return; + expect(d.next.clearSince).toBeNull(); + expect(d.next.sends).toBe(1); + expect(d.next.since).toBe(0); + }); + + it("a sustained clear then fresh breach restarts at send #1", () => { + const state = firstSend(0); + const started = clear(state, 6 * MIN); + if (started.kind !== "commit" || !started.next) throw new Error("expected grace start"); + const recovered = clear(started.next, 6 * MIN + P.recoveryGraceMs); + if (recovered.kind !== "commit" || recovered.next !== null) throw new Error("expected recovery"); + const fresh = breach(null, 10 * MIN); + expect(fresh.kind).toBe("send"); + if (fresh.kind !== "send") return; + expect(fresh.next.sends).toBe(1); + expect(fresh.next.since).toBe(10 * MIN); + }); +}); + +describe("commitAfterSend", () => { + const decision = () => { + const d = breach(null, 0); + if (d.kind !== "send") throw new Error("expected send"); + return d; + }; + + it("persists the new incident only on a confirmed send", () => { + expect(commitAfterSend(null, decision(), { status: "sent" })).toEqual(decision().next); + }); + + it("keeps the prior state when the send failed or was skipped", () => { + const prior = firstSend(0); + expect(commitAfterSend(prior, decision(), { status: "failed", error: "x" })).toBe(prior); + expect(commitAfterSend(prior, decision(), { status: "skipped", reason: "no_recipient" })).toBe(prior); + expect(commitAfterSend(null, decision(), { status: "failed", error: "x" })).toBeNull(); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/incident-tracker.test.ts b/apps/api/src/monitoring/__tests__/incident-tracker.test.ts new file mode 100644 index 0000000..59b0e9f --- /dev/null +++ b/apps/api/src/monitoring/__tests__/incident-tracker.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, it } from "bun:test"; +import type { Incident } from "../incident-policy"; +import { createIncidentTracker, type IncidentStore } from "../incident-tracker"; + +const MIN = 60_000; + +const makeStore = () => { + const data = new Map(); + let loadError = false; + const store: IncidentStore = { + async load(id) { + if (loadError) throw new Error("redis down"); + return data.get(id) ?? null; + }, + async save(id, state) { + if (state) data.set(id, state); + else data.delete(id); + }, + }; + return { store, data, failLoad: () => (loadError = true) }; +}; + +const makeSend = () => { + const calls = { count: 0 }; + let failuresLeft = 0; + const send = async () => { + calls.count++; + if (failuresLeft > 0) { + failuresLeft--; + return { status: "failed", error: "smtp down" } as const; + } + return { status: "sent" } as const; + }; + return { calls, send, failNext: () => (failuresLeft = 1) }; +}; + +describe("incident tracker", () => { + it("sends the first breach and stores the incident", async () => { + const { store, data } = makeStore(); + const { calls, send } = makeSend(); + const tracker = createIncidentTracker(store); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send }); + + expect(calls.count).toBe(1); + expect(data.get("a1")?.sends).toBe(1); + }); + + it("suppresses repeat sends until the backoff is due", async () => { + const { store, data } = makeStore(); + const { calls, send } = makeSend(); + const tracker = createIncidentTracker(store); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send }); + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 4 * MIN, send }); + expect(calls.count).toBe(1); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 5 * MIN, send }); + expect(calls.count).toBe(2); + expect(data.get("a1")?.sends).toBe(2); + }); + + it("keeps the prior state when the send fails, then retries next tick", async () => { + const { store, data } = makeStore(); + const { calls, send, failNext } = makeSend(); + const tracker = createIncidentTracker(store); + + failNext(); + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send }); + expect(calls.count).toBe(1); + expect(data.has("a1")).toBe(false); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: MIN, send }); + expect(calls.count).toBe(2); + expect(data.get("a1")?.sends).toBe(1); + }); + + it("deletes the incident after a sustained clear", async () => { + const { store, data } = makeStore(); + const { send } = makeSend(); + const tracker = createIncidentTracker(store); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send }); + await tracker.step({ alertId: "a1", observation: { kind: "clear" }, now: 2 * MIN, send }); + expect(data.get("a1")?.clearSince).toBe(2 * MIN); + + await tracker.step({ alertId: "a1", observation: { kind: "clear" }, now: 8 * MIN, send }); + expect(data.has("a1")).toBe(false); + }); + + it("refreshes the stored incident on a suppressed tick", async () => { + const { store, data } = makeStore(); + const { calls, send } = makeSend(); + const tracker = createIncidentTracker(store); + + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send }); + const stored = data.get("a1"); + await tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: MIN, send }); + expect(calls.count).toBe(1); + expect(data.get("a1")).toEqual(stored); + }); + + it("aborts before sending when the store is unavailable", async () => { + const { store, failLoad } = makeStore(); + const { calls, send } = makeSend(); + const tracker = createIncidentTracker(store); + + failLoad(); + await expect(tracker.step({ alertId: "a1", observation: { kind: "breach" }, now: 0, send })).rejects.toThrow( + "redis down", + ); + expect(calls.count).toBe(0); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/slack-message.test.ts b/apps/api/src/monitoring/__tests__/slack-message.test.ts new file mode 100644 index 0000000..f985b39 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/slack-message.test.ts @@ -0,0 +1,88 @@ +import { describe, expect, test } from "bun:test"; +import { buildSlackMessage } from "../slack-message"; + +describe("buildSlackMessage", () => { + const findBlock = (blocks: any[], type: string) => blocks.find((b) => b.type === type); + + test("cpu alert formats threshold and value as percentages", () => { + const { text, blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6); + expect(text).toBe("tabi: cpu alert"); + const fields = findBlock(blocks, "section").fields.map((f: any) => f.text); + expect(fields).toContain("*Type:* cpu"); + expect(fields).toContain("*Threshold:* 80%"); + expect(fields).toContain("*Current:* 51.6%"); + }); + + test("memory alert shows MB for current value", () => { + const { blocks } = buildSlackMessage("memory", "tabi", 85, 524.4); + const fields = findBlock(blocks, "section").fields.map((f: any) => f.text); + expect(fields).toContain("*Current:* 524 MB"); + expect(fields).toContain("*Threshold:* 85"); + }); + + test("downtime alert shows service down", () => { + const { blocks } = buildSlackMessage("downtime", "tabi", null, 1); + const fields = findBlock(blocks, "section").fields.map((f: any) => f.text); + expect(fields).toContain("*Type:* downtime"); + expect(fields).toContain("*Threshold:* N/A"); + expect(fields).toContain("*Current:* Service down"); + }); + + test("includes per-container values when details present", () => { + const { blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6, { + containers: [ + { name: "tabi-abc", value: 51.6 }, + { name: "tabi-def", value: 49.2 }, + ], + }); + const section = blocks.find((b: any) => b.type === "section" && b.text?.text?.includes("tabi-abc")); + expect(section).toBeDefined(); + expect(section.text.text).toContain("`tabi-abc` 51.6%"); + expect(section.text.text).toContain("`tabi-def` 49.2%"); + }); + + test("scaling suggestion enable_autoscaling renders title and button", () => { + const { blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6, { + scaling: { kind: "enable_autoscaling", url: "https://dequel.example/project/1?tab=scaling" }, + }); + const section = blocks.find((b: any) => b.type === "section" && b.text?.text?.includes("Autoscaling is off")); + expect(section).toBeDefined(); + expect(section.text.text).toContain("Enable autoscaling"); + const actions = blocks.filter((b: any) => b.type === "actions"); + const cta = actions.flatMap((a: any) => a.elements).find((e: any) => e.text.text === "Set up autoscaling"); + expect(cta).toBeDefined(); + expect(cta.url).toBe("https://dequel.example/project/1?tab=scaling"); + }); + + test("scaling suggestion increase_max_replicas shows replica limit and counts", () => { + const { blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6, { + scaling: { kind: "increase_max_replicas", current: 3, maxReplicas: 3, url: "https://x/scaling" }, + }); + const section = blocks.find((b: any) => b.text?.text?.includes("Replica limit reached")); + expect(section.text.text).toContain("(3/3)"); + const cta = blocks + .filter((b: any) => b.type === "actions") + .flatMap((a: any) => a.elements) + .find((e: any) => e.text.text === "Adjust scaling"); + expect(cta).toBeDefined(); + }); + + test("projectUrl renders Open project button", () => { + const { blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6, { + projectUrl: "https://dequel.example/project/42", + }); + const btn = blocks + .filter((b: any) => b.type === "actions") + .flatMap((a: any) => a.elements) + .find((e: any) => e.text.text === "Open project"); + expect(btn).toBeDefined(); + expect(btn.url).toBe("https://dequel.example/project/42"); + }); + + test("no scaling block and no project button without details", () => { + const { blocks } = buildSlackMessage("cpu", "tabi", 80, 51.6); + expect(blocks.filter((b: any) => b.type === "actions")).toHaveLength(0); + expect(blocks.some((b: any) => b.text?.text?.includes("Autoscaling"))).toBe(false); + expect(blocks).toHaveLength(2); + }); +}); diff --git a/apps/api/src/monitoring/__tests__/templates.test.ts b/apps/api/src/monitoring/__tests__/templates.test.ts new file mode 100644 index 0000000..9e56d69 --- /dev/null +++ b/apps/api/src/monitoring/__tests__/templates.test.ts @@ -0,0 +1,138 @@ +import { describe, expect, it } from "bun:test"; +import { buildDeploymentFailureEmail, buildEmail, buildSmtpTestEmail } from "../templates"; + +describe("Email Templates", () => { + describe("buildDeploymentFailureEmail", () => { + it("generates a deployment failure email matching Dequel design system", () => { + const { subject, html } = buildDeploymentFailureEmail({ + projectName: "clinsight-be", + failureReason: "fix(docker): add build dependencies for pycairo compilation (#29) (#30)", + commitSha: "a1b2c3d4e5f67890", + sourceRef: "main", + finishedAt: "2026-09-24T01:43:00Z", + logsUrl: "https://dequel.app/project/proj-123?tab=deployments", + }); + + expect(subject).toBe("deploy failed for clinsight-be"); + expect(html).toContain("Dequel"); + expect(html).toContain("DEPLOY FAILED"); + expect(html).toContain("clinsight-be"); + expect(html).toContain("a1b2c3d4e5f6"); + expect(html).toContain("main"); + expect(html).toContain("fix(docker): add build dependencies"); + expect(html).toContain("View Logs"); + expect(html).toContain("https://dequel.app/project/proj-123?tab=deployments"); + expect(html).toContain("Learn more"); + expect(html).toContain("troubleshooting deploys on"); + expect(html).toContain('The Dequel team'); + // Check WebP Logo & HTML table Grid graphic presence + expect(html).toContain("logo_xzvwej.webp"); + expect(html).toContain("border-spacing:3px"); + }); + + it("handles missing optional context fields gracefully", () => { + const { subject, html } = buildDeploymentFailureEmail({ + projectName: "api-service", + failureReason: null, + commitSha: null, + sourceRef: "dev", + finishedAt: null, + }); + + expect(subject).toBe("deploy failed for api-service"); + expect(html).toContain("DEPLOY FAILED"); + expect(html).not.toContain("View Logs"); + }); + }); + + describe("buildEmail (monitoring alerts)", () => { + it("generates CPU alert email", () => { + const { subject, html } = buildEmail("cpu", "web-app", 80, 94.2, { + containers: [{ name: "web-app-1", value: 94.2 }], + appUrl: "https://web-app.example.com", + }); + + expect(subject).toBe("High CPU on web-app (94%)"); + expect(html).toContain("CPU ALERT"); + expect(html).toContain("94.2%"); + expect(html).toContain("View Application"); + }); + + it("generates Memory alert email", () => { + const { subject, html } = buildEmail("memory", "worker-app", 85, 91.0); + + expect(subject).toBe("High memory on worker-app (91%)"); + expect(html).toContain("MEMORY ALERT"); + expect(html).toContain("91.0%"); + }); + + it("renders the enable-autoscaling suggestion", () => { + const { html } = buildEmail("cpu", "web-app", 80, 94.2, { + scaling: { kind: "enable_autoscaling", url: "https://dequel.local/project/p1?tab=scaling" }, + }); + + expect(html).toContain("Autoscaling is off"); + expect(html).toContain("Set up autoscaling"); + expect(html).not.toContain("SCALING_HTML"); + }); + + it("renders the increase-replicas suggestion", () => { + const { html } = buildEmail("memory", "web-app", 85, 91.0, { + scaling: { + kind: "increase_max_replicas", + current: 3, + maxReplicas: 3, + url: "https://dequel.local/project/p1?tab=scaling", + }, + }); + + expect(html).toContain("Replica limit reached"); + expect(html).toContain("(3/3)"); + expect(html).toContain("Adjust scaling"); + }); + + it("renders no scaling box without a suggestion", () => { + const { html } = buildEmail("cpu", "web-app", 80, 94.2); + expect(html).not.toContain("SCALING_HTML"); + expect(html).not.toContain("Autoscaling is off"); + }); + + it("generates downtime alert email", () => { + const { subject, html } = buildEmail("downtime", "db-proxy", null, 0, { + lastRunningAt: "2026-09-26T08:00:00Z", + logsUrl: "https://dequel.app/logs", + }); + + expect(subject).toBe("db-proxy is down"); + expect(html).toContain("SERVICE DOWN"); + expect(html).toContain("Offline"); + expect(html).toContain("View Logs"); + }); + + it("generates cert expiry alert email", () => { + const { subject, html } = buildEmail("cert_expiry", "my-domain.com", null, 14); + + expect(subject).toBe("SSL certificate for my-domain.com expires in 14 days"); + expect(html).toContain("CERTIFICATE EXPIRY"); + expect(html).toContain("14"); + }); + + it("generates default fallback alert email", () => { + const { subject, html } = buildEmail("disk_space", "storage-node", 90, 95); + + expect(subject).toBe("[Dequel] disk_space alert — storage-node"); + expect(html).toContain("ALERT"); + }); + }); + + describe("buildSmtpTestEmail", () => { + it("generates a formatted SMTP test email", () => { + const { subject, html } = buildSmtpTestEmail(); + + expect(subject).toBe("[Dequel] SMTP Test Email"); + expect(html).toContain("SMTP TEST"); + expect(html).toContain("SMTP Transport Verified"); + expect(html).toContain("Dequel"); + }); + }); +}); diff --git a/apps/api/src/monitoring/alert-guard.ts b/apps/api/src/monitoring/alert-guard.ts new file mode 100644 index 0000000..fb6338d --- /dev/null +++ b/apps/api/src/monitoring/alert-guard.ts @@ -0,0 +1,30 @@ +export type ScalingSuggestion = + | { kind: "enable_autoscaling" } + | { kind: "increase_max_replicas"; current: number; maxReplicas: number }; + +export type ScalingContext = { + policy: { enabled: boolean; maxReplicas: number } | null; + cpuLimit: number | null; + currentReplicas: number | null; +}; + +export type ScalingGuard = + | { suppress: true; suggestion: null } + | { suppress: false; suggestion: ScalingSuggestion | null }; + +export const scalingGuard = (alertType: string, ctx: ScalingContext): ScalingGuard => { + if (alertType !== "cpu" && alertType !== "memory") { + return { suppress: false, suggestion: null }; + } + if (!ctx.policy || !ctx.policy.enabled || !ctx.cpuLimit || ctx.cpuLimit <= 0) { + return { suppress: false, suggestion: { kind: "enable_autoscaling" } }; + } + const current = ctx.currentReplicas ?? 1; + if (current >= ctx.policy.maxReplicas) { + return { + suppress: false, + suggestion: { kind: "increase_max_replicas", current, maxReplicas: ctx.policy.maxReplicas }, + }; + } + return { suppress: true, suggestion: null }; +}; diff --git a/apps/api/src/monitoring/container-stats.ts b/apps/api/src/monitoring/container-stats.ts new file mode 100644 index 0000000..5f1c164 --- /dev/null +++ b/apps/api/src/monitoring/container-stats.ts @@ -0,0 +1,84 @@ +import { getServerById } from "../db/repo"; +import { run } from "../orchestrator/runtime"; +import type { Deployment, Server } from "../types"; +import { dockerBin } from "../utils/docker-bin"; + +const AGENT_OFFLINE_MS = 90_000; + +export interface ContainerStats { + cpuPercent: number; + memoryMb: number; +} + +const parseMemToMb = (mem: string): number => { + const match = mem.match(/^([\d.]+)(\w+)$/); + if (!match) return 0; + const val = parseFloat(match[1]); + switch (match[2]) { + case "GiB": + case "GB": + return val * 1024; + case "MiB": + case "MB": + return val; + case "KiB": + case "KB": + return val / 1024; + default: + return val; + } +}; + +const parseStatsJson = (statsJson: string): ContainerStats | null => { + try { + const stats = JSON.parse(statsJson); + return { + cpuPercent: parseFloat(stats.CPUPerc?.replace("%", "") ?? "0"), + memoryMb: parseMemToMb(stats.MemUsage?.split("/")[0]?.trim() ?? "0B"), + }; + } catch { + return null; + } +}; + +const isAgentOffline = (server: Server | null): boolean => { + if (!server?.lastHeartbeat) return true; + return Date.now() - new Date(server.lastHeartbeat).getTime() > AGENT_OFFLINE_MS; +}; + +export const getDeploymentContainerStats = async (deployment: Deployment): Promise => { + const server = + deployment.serverId && deployment.serverId !== "local" + ? await getServerById(deployment.serverId).catch(() => null) + : null; + const mode = server?.mode ?? "local"; + if (mode === "agent") { + if (isAgentOffline(server)) return null; + const { agentStatsCache } = await import("../agents/stats-cache"); + const containers = await agentStatsCache.get(server!.id); + const stat = containers.get(deployment.containerName ?? ""); + return stat ? { cpuPercent: stat.cpuPercent, memoryMb: stat.memoryMb } : null; + } + try { + const statsJson = await run( + dockerBin, + ["stats", "--no-stream", "--format", "{{json .}}", deployment.containerName ?? ""], + server, + ); + return parseStatsJson(statsJson); + } catch { + return null; + } +}; + +export const collectContainerStats = async ( + deployments: Deployment[], +): Promise<{ name: string; cpuPercent: number; memoryMb: number }[]> => { + const out: { name: string; cpuPercent: number; memoryMb: number }[] = []; + for (const dep of deployments) { + if (dep.status !== "running" || !dep.containerName) continue; + const stats = await getDeploymentContainerStats(dep); + if (stats) out.push({ name: dep.containerName, cpuPercent: stats.cpuPercent, memoryMb: stats.memoryMb }); + } + return out; +}; diff --git a/apps/api/src/monitoring/evaluator.ts b/apps/api/src/monitoring/evaluator.ts index b756452..191bfa0 100644 --- a/apps/api/src/monitoring/evaluator.ts +++ b/apps/api/src/monitoring/evaluator.ts @@ -1,115 +1,28 @@ import { eq } from "drizzle-orm"; import Redis from "ioredis"; import { getDb } from "../db/db-provider"; -import { getProjectById, getServerById, listDeployments } from "../db/repo"; +import { getProjectById, getScalingPolicy, listDeployments } from "../db/repo"; import { alerts } from "../db/schema"; -import { run } from "../orchestrator/runtime"; -import type { Deployment, Server } from "../types"; +import { scalingEngine } from "../scaling/engine"; +import type { Deployment } from "../types"; import { config } from "../utils/config"; -import { dockerBin } from "../utils/docker-bin"; +import { appBaseUrl } from "../utils/routes"; +import { scalingGuard } from "./alert-guard"; +import { getDeploymentContainerStats } from "./container-stats"; +import { type Observation } from "./incident-policy"; +import { createIncidentTracker, createRedisIncidentStore } from "./incident-tracker"; import { sendNotification } from "./notifier"; - -const NOTIFICATION_KEY = "dequel:alert:notified"; -const NOTIFICATION_COOLDOWN_MS = 300_000; // 5 min between same alert -const AGENT_OFFLINE_MS = 90_000; - -interface ContainerStats { - cpuPercent: number; - memoryMb: number; -} - -const parseMemToMb = (mem: string): number => { - const match = mem.match(/^([\d.]+)(\w+)$/); - if (!match) return 0; - const val = parseFloat(match[1]); - switch (match[2]) { - case "GiB": - case "GB": - return val * 1024; - case "MiB": - case "MB": - return val; - case "KiB": - case "KB": - return val / 1024; - default: - return val; - } -}; - -const parseStatsJson = (statsJson: string): ContainerStats | null => { - try { - const stats = JSON.parse(statsJson); - return { - cpuPercent: parseFloat(stats.CPUPerc?.replace("%", "") ?? "0"), - memoryMb: parseMemToMb(stats.MemUsage?.split("/")[0]?.trim() ?? "0B"), - }; - } catch { - return null; - } -}; - -const isAgentOffline = (server: Server | null): boolean => { - if (!server?.lastHeartbeat) return true; - return Date.now() - new Date(server.lastHeartbeat).getTime() > AGENT_OFFLINE_MS; -}; - -const getContainerStats = async (deployment: Deployment): Promise => { - const server = - deployment.serverId && deployment.serverId !== "local" - ? await getServerById(deployment.serverId).catch(() => null) - : null; - const mode = server?.mode ?? "local"; - if (mode === "agent") { - if (isAgentOffline(server)) return null; - const { agentStatsCache } = await import("../agents/stats-cache"); - const containers = await agentStatsCache.get(server!.id); - const stat = containers.get(deployment.containerName ?? ""); - return stat ? { cpuPercent: stat.cpuPercent, memoryMb: stat.memoryMb } : null; - } - try { - const statsJson = await run( - dockerBin, - ["stats", "--no-stream", "--format", "{{json .}}", deployment.containerName ?? ""], - server, - ); - return parseStatsJson(statsJson); - } catch { - return null; - } -}; - -const getMetricValue = async (alertType: string, _projectId: string, deployments: Deployment[]): Promise => { - if (alertType === "cpu" || alertType === "memory") { - let total = 0; - let count = 0; - for (const dep of deployments) { - if (dep.status !== "running" || !dep.containerName) continue; - const stats = await getContainerStats(dep); - if (stats) { - total += alertType === "cpu" ? stats.cpuPercent : stats.memoryMb; - count++; - } - } - return count > 0 ? total / count : 0; - } - if (alertType === "downtime") { - const running = deployments.filter((d) => d.status === "running"); - return running.length === 0 ? 1 : 0; - } - if (alertType === "error_rate") { - const failed = deployments.filter((d) => d.status === "failed"); - return failed.length > 0 ? failed.length : 0; - } - return 0; -}; +import type { AlertDetails } from "./templates"; class AlertEvaluator { private redis: Redis; + private incidents: ReturnType; private interval: ReturnType | null = null; + private ticking = false; constructor() { this.redis = new Redis(config.redisUrl, { maxRetriesPerRequest: null, enableOfflineQueue: false }); + this.incidents = createIncidentTracker(createRedisIncidentStore(this.redis)); } start() { @@ -128,6 +41,8 @@ class AlertEvaluator { } private async tick() { + if (this.ticking) return; + this.ticking = true; try { const db = await getDb(); const alertRows = await db.select().from(alerts).where(eq(alerts.enabled, true)).execute(); @@ -156,52 +71,128 @@ class AlertEvaluator { } } catch (err) { console.error("[Alerts] Tick error:", err); + } finally { + this.ticking = false; } } - private async evaluate(alert: any, project: { id: string; name: string }, deployments: Deployment[]) { - const currentValue = await getMetricValue(alert.type, project.id, deployments); - if (currentValue === 0) return; - + private async evaluate( + alert: any, + project: { id: string; name: string; liveUrl?: string | null; cpuLimit?: number | null }, + deployments: Deployment[], + ) { const threshold = alert.threshold ?? (alert.type === "memory" ? 85 : 70); - let breached = false; - - switch (alert.type) { - case "cpu": - breached = currentValue > threshold; - break; - case "memory": - breached = currentValue > threshold; - break; - case "downtime": - breached = currentValue > 0; - break; - case "error_rate": - breached = currentValue > 0; - break; - case "cert_expiry": - break; + let { observation, currentValue } = await this.probe(alert.type, deployments, threshold); + let scaling: AlertDetails["scaling"]; + + if (observation.kind === "breach") { + const guard = await this.evaluateScalingGuard(alert.type, project); + if (guard.suppress) { + observation = { kind: "no_data" }; + } else if (guard.suggestion) { + scaling = { + ...guard.suggestion, + url: `${appBaseUrl()}/project/${project.id}?tab=scaling`, + }; + } } - if (!breached) return; - - const notifiedKey = `${NOTIFICATION_KEY}:${alert.id}`; - const lastNotified = await this.redis - .get(notifiedKey) - .then((v) => (v ? Number(v) : 0)) - .catch(() => 0); - if (Date.now() - lastNotified < NOTIFICATION_COOLDOWN_MS) return; - - await sendNotification({ - channel: alert.channel, - destination: alert.destination, - projectName: project.name, - alertType: alert.type, - threshold, - currentValue, + await this.incidents.step({ + alertId: alert.id, + observation, + send: async () => { + const details = await this.buildDetails(alert.type, deployments); + return sendNotification({ + channel: alert.channel, + destination: alert.destination, + projectName: project.name, + alertType: alert.type, + threshold, + currentValue, + details: { + ...details, + scaling, + projectUrl: `${appBaseUrl()}/project/${project.id}`, + }, + }); + }, + }); + } + + private async evaluateScalingGuard(alertType: string, project: { id: string; cpuLimit?: number | null }) { + const policy = await getScalingPolicy(project.id).catch(() => null); + const canScale = !!policy?.enabled && !!project.cpuLimit && project.cpuLimit > 0; + const replicas = canScale ? await scalingEngine.getProjectReplicas(project.id) : null; + return scalingGuard(alertType, { + policy: policy ? { enabled: policy.enabled, maxReplicas: policy.maxReplicas } : null, + cpuLimit: project.cpuLimit ?? null, + currentReplicas: replicas?.current ?? null, }); + } + + private async probe( + alertType: string, + deployments: Deployment[], + threshold: number, + ): Promise<{ observation: Observation; currentValue: number }> { + if (alertType === "downtime") { + if (!deployments.length) return { observation: { kind: "no_data" }, currentValue: 0 }; + const running = deployments.some((d) => d.status === "running"); + return { + observation: { kind: running ? "clear" : "breach" }, + currentValue: running ? 0 : 1, + }; + } + + if (alertType !== "cpu" && alertType !== "memory") { + return { observation: { kind: "no_data" }, currentValue: 0 }; + } + + let total = 0; + let count = 0; + for (const dep of deployments) { + if (dep.status !== "running" || !dep.containerName) continue; + const stats = await getDeploymentContainerStats(dep); + if (stats) { + total += alertType === "cpu" ? stats.cpuPercent : stats.memoryMb; + count++; + } + } + if (count === 0) return { observation: { kind: "no_data" }, currentValue: 0 }; + + const value = total / count; + return { + observation: { kind: value > threshold ? "breach" : "clear" }, + currentValue: value, + }; + } + + private async buildDetails(alertType: string, deployments: Deployment[]): Promise { + const details: AlertDetails = {}; + + if (alertType === "cpu" || alertType === "memory") { + const containers: { name: string; value: number }[] = []; + for (const dep of deployments) { + if (dep.status !== "running" || !dep.containerName) continue; + const stats = await getDeploymentContainerStats(dep); + if (stats) { + containers.push({ + name: dep.containerName, + value: alertType === "cpu" ? stats.cpuPercent : stats.memoryMb, + }); + } + } + details.containers = containers; + } + + if (alertType === "downtime") { + const lastFinished = deployments + .filter((d) => d.finishedAt) + .sort((a, b) => new Date(b.finishedAt!).getTime() - new Date(a.finishedAt!).getTime())[0]; + details.lastRunningAt = lastFinished?.finishedAt ?? null; + } - await this.redis.set(notifiedKey, String(Date.now())).catch(() => {}); + return details; } } diff --git a/apps/api/src/monitoring/failure-notifier.ts b/apps/api/src/monitoring/failure-notifier.ts new file mode 100644 index 0000000..19388c9 --- /dev/null +++ b/apps/api/src/monitoring/failure-notifier.ts @@ -0,0 +1,85 @@ +import { + claimFailureNotification, + listPendingFailureNotificationIds, + markFailureNotificationSent, +} from "../db/repo/deployment-events"; +import { onDeploymentFailed } from "../events"; +import { config } from "../utils/config"; +import { isSmtpConfigured, sendDeploymentFailureEmail } from "./notifier"; + +let queue: string[] = []; +const pending = new Set(); +let draining = false; +let sweepTimer: ReturnType | null = null; +let unsubscribe: (() => void) | null = null; +let configuredWarned = false; + +const enqueue = (eventId: string) => { + if (pending.has(eventId)) return; + pending.add(eventId); + queue.push(eventId); + void drain(); +}; + +const drain = async () => { + if (draining) return; + draining = true; + try { + while (queue.length > 0) { + const eventId = queue.shift()!; + pending.delete(eventId); + try { + if (!(await isSmtpConfigured())) { + if (!configuredWarned) { + console.warn("[FailureNotifier] SMTP not configured — leaving failure events pending"); + configuredWarned = true; + } + continue; + } + configuredWarned = false; + const ctx = await claimFailureNotification(eventId); + if (!ctx) continue; + const delivery = await sendDeploymentFailureEmail(ctx); + if (delivery.status === "sent") { + await markFailureNotificationSent(eventId); + console.log(`[FailureNotifier] failure email sent for ${ctx.deploymentId} (attempt ${ctx.attempt})`); + } else { + console.warn( + `[FailureNotifier] ${ctx.deploymentId}: ${delivery.status}`, + "reason" in delivery ? delivery.reason : delivery.error, + ); + } + } catch (err) { + console.error(`[FailureNotifier] ${eventId}:`, err); + } + } + } finally { + draining = false; + } +}; + +const sweep = async () => { + try { + const ids = await listPendingFailureNotificationIds(); + for (const id of ids) enqueue(id); + } catch (err) { + console.error("[FailureNotifier] sweep failed:", err); + } +}; + +export const startFailureNotifier = (): void => { + if (unsubscribe) return; + unsubscribe = onDeploymentFailed((signal) => enqueue(signal.eventId)); + sweepTimer = setInterval(() => void sweep(), config.failureSweepIntervalMs); + void sweep(); + console.log("[FailureNotifier] started"); +}; + +export const stopFailureNotifier = (): void => { + unsubscribe?.(); + unsubscribe = null; + if (sweepTimer) clearInterval(sweepTimer); + sweepTimer = null; + queue = []; + pending.clear(); +}; diff --git a/apps/api/src/monitoring/health.ts b/apps/api/src/monitoring/health.ts new file mode 100644 index 0000000..59a19e1 --- /dev/null +++ b/apps/api/src/monitoring/health.ts @@ -0,0 +1,66 @@ +import type { HealthCheck, OverallHealth } from "../types"; + +const SERVER_HEARTBEAT_FAIL_MS = 90_000; +const HTTP_ERROR_RATE_WARN = 0.05; + +export const evaluateHealth = (input: { + server: { status: string; lastHeartbeatAgeMs: number | null } | null; + runningDeployments: number; + hasDeployments: boolean; + routeErrors: string[]; + http: { available: boolean; errorRate: number | null }; +}): { overall: OverallHealth; checks: HealthCheck[] } => { + const checks: HealthCheck[] = []; + + if (!input.server) { + checks.push({ name: "server", status: "ok", detail: null }); + } else if (input.server.lastHeartbeatAgeMs === null) { + checks.push({ name: "server", status: "fail", detail: "Server never heartbeated" }); + } else if (input.server.lastHeartbeatAgeMs > SERVER_HEARTBEAT_FAIL_MS) { + checks.push({ + name: "server", + status: "fail", + detail: `Heartbeat stale for ${Math.round(input.server.lastHeartbeatAgeMs / 1000)}s`, + }); + } else { + checks.push({ name: "server", status: "ok", detail: null }); + } + + if (!input.hasDeployments) { + checks.push({ name: "containers", status: "warn", detail: "No deployments yet" }); + } else if (input.runningDeployments === 0) { + checks.push({ name: "containers", status: "fail", detail: "No running containers" }); + } else { + checks.push({ + name: "containers", + status: "ok", + detail: `${input.runningDeployments} running`, + }); + } + + if (input.routeErrors.length > 0) { + checks.push({ name: "ingress", status: "warn", detail: input.routeErrors[0] }); + } else { + checks.push({ name: "ingress", status: "ok", detail: null }); + } + + if (!input.http.available) { + checks.push({ name: "http", status: "unknown", detail: "Loki unavailable" }); + } else if (input.http.errorRate === null) { + checks.push({ name: "http", status: "unknown", detail: "No traffic" }); + } else if (input.http.errorRate > HTTP_ERROR_RATE_WARN) { + checks.push({ + name: "http", + status: "warn", + detail: `${(input.http.errorRate * 100).toFixed(1)}% error rate`, + }); + } else { + checks.push({ name: "http", status: "ok", detail: null }); + } + + let overall: OverallHealth = "healthy"; + if (checks.some((c) => c.status === "fail")) overall = "down"; + else if (checks.some((c) => c.status === "warn" || c.status === "unknown")) overall = "degraded"; + + return { overall, checks }; +}; diff --git a/apps/api/src/monitoring/http-error-rate.ts b/apps/api/src/monitoring/http-error-rate.ts new file mode 100644 index 0000000..a2c6364 --- /dev/null +++ b/apps/api/src/monitoring/http-error-rate.ts @@ -0,0 +1,46 @@ +import type { HttpErrorRate } from "../types"; +import { buildProjectRequestHostRegex, caddyRequestLogSelector, lokiInstantQuery } from "../utils/loki"; + +const unavailable = (windowSeconds: number): HttpErrorRate => ({ + source: "loki", + available: false, + windowSeconds, + totalRequests: 0, + errorRequests: 0, + errorRate: null, + byStatus: [], +}); + +export const parseHttpErrorRate = (result: unknown, windowSeconds: number): HttpErrorRate => { + if (!Array.isArray(result)) return unavailable(windowSeconds); + const byStatus: { status: string; count: number }[] = []; + for (const entry of result as any[]) { + const status = String(entry?.metric?.status ?? ""); + const count = Number(entry?.value?.[1]); + if (!status || !Number.isFinite(count)) continue; + byStatus.push({ status, count }); + } + const totalRequests = byStatus.reduce((sum, b) => sum + b.count, 0); + const errorRequests = byStatus.reduce((sum, b) => sum + (Number(b.status) >= 400 ? b.count : 0), 0); + return { + source: "loki", + available: true, + windowSeconds, + totalRequests, + errorRequests, + errorRate: totalRequests > 0 ? errorRequests / totalRequests : null, + byStatus, + }; +}; + +export const getHttpErrorRate = async (projectId: string, windowSeconds: number): Promise => { + try { + const hostRegex = await buildProjectRequestHostRegex(projectId); + const query = `sum by (status) (count_over_time(${caddyRequestLogSelector(hostRegex)} [${windowSeconds}s]))`; + const result = await lokiInstantQuery(query); + if (result === null) return unavailable(windowSeconds); + return parseHttpErrorRate(result, windowSeconds); + } catch { + return unavailable(windowSeconds); + } +}; diff --git a/apps/api/src/monitoring/incident-policy.ts b/apps/api/src/monitoring/incident-policy.ts new file mode 100644 index 0000000..9f71255 --- /dev/null +++ b/apps/api/src/monitoring/incident-policy.ts @@ -0,0 +1,74 @@ +import type { MailDelivery } from "../types"; + +export type Observation = { kind: "breach" } | { kind: "clear" } | { kind: "no_data" }; + +export type Incident = { + status: "active"; + since: number; + sends: number; + nextDueAt: number; + clearSince: number | null; +}; + +export type IncidentPolicy = { + baseMs: number; + factor: number; + capMs: number; + recoveryGraceMs: number; +}; + +export type Decision = { kind: "noop" } | { kind: "commit"; next: Incident | null } | { kind: "send"; next: Incident }; + +export const DEFAULT_POLICY: IncidentPolicy = { + baseMs: 300_000, + factor: 4, + capMs: 43_200_000, + recoveryGraceMs: 300_000, +}; + +export const backoffMs = (sends: number, p: IncidentPolicy): number => { + const exponent = Math.max(0, sends - 1); + if (exponent > 30) return p.capMs; + return Math.min(p.baseMs * p.factor ** exponent, p.capMs); +}; + +export const decide = (state: Incident | null, obs: Observation, now: number, p: IncidentPolicy): Decision => { + if (obs.kind === "no_data") return { kind: "noop" }; + + if (obs.kind === "clear") { + if (!state) return { kind: "noop" }; + if (state.clearSince === null) return { kind: "commit", next: { ...state, clearSince: now } }; + if (now - state.clearSince >= p.recoveryGraceMs) return { kind: "commit", next: null }; + return { kind: "noop" }; + } + + if (!state) { + return { + kind: "send", + next: { + status: "active", + since: now, + sends: 1, + nextDueAt: now + backoffMs(1, p), + clearSince: null, + }, + }; + } + + const reentered = state.clearSince !== null ? { ...state, clearSince: null } : state; + if (now >= reentered.nextDueAt) { + const sends = reentered.sends + 1; + return { + kind: "send", + next: { ...reentered, sends, nextDueAt: now + backoffMs(sends, p) }, + }; + } + if (reentered !== state) return { kind: "commit", next: reentered }; + return { kind: "noop" }; +}; + +export const commitAfterSend = ( + prior: Incident | null, + decision: Extract, + delivery: MailDelivery, +): Incident | null => (delivery.status === "sent" ? decision.next : prior); diff --git a/apps/api/src/monitoring/incident-tracker.ts b/apps/api/src/monitoring/incident-tracker.ts new file mode 100644 index 0000000..b7bc617 --- /dev/null +++ b/apps/api/src/monitoring/incident-tracker.ts @@ -0,0 +1,82 @@ +import type { Redis } from "ioredis"; +import type { MailDelivery } from "../types"; +import { + type Decision, + type Incident, + type IncidentPolicy, + type Observation, + DEFAULT_POLICY, + commitAfterSend, + decide, +} from "./incident-policy"; + +const KEY_PREFIX = "dequel:alert:incident:"; +const TTL_SECONDS = 604_800; + +export type StepInput = { + alertId: string; + observation: Observation; + now?: number; + send: () => Promise; +}; + +export type IncidentStore = { + load(alertId: string): Promise; + save(alertId: string, state: Incident | null): Promise; +}; + +export type IncidentTracker = { step(input: StepInput): Promise }; + +const isIncident = (value: unknown): value is Incident => { + if (!value || typeof value !== "object") return false; + const v = value as Record; + return ( + v.status === "active" && + typeof v.since === "number" && + typeof v.sends === "number" && + typeof v.nextDueAt === "number" && + (v.clearSince === null || typeof v.clearSince === "number") + ); +}; + +export const createRedisIncidentStore = (redis: Redis): IncidentStore => ({ + async load(alertId) { + const raw = await redis.get(KEY_PREFIX + alertId); + if (!raw) return null; + try { + const parsed = JSON.parse(raw); + return isIncident(parsed) ? parsed : null; + } catch { + return null; + } + }, + async save(alertId, state) { + const key = KEY_PREFIX + alertId; + if (!state) { + await redis.del(key); + return; + } + await redis.set(key, JSON.stringify(state), "EX", TTL_SECONDS); + }, +}); + +export const createIncidentTracker = ( + store: IncidentStore, + policy: IncidentPolicy = DEFAULT_POLICY, +): IncidentTracker => ({ + async step(input) { + const state = await store.load(input.alertId); + const decision: Decision = decide(state, input.observation, input.now ?? Date.now(), policy); + + if (decision.kind === "send") { + const delivery = await input.send(); + await store.save(input.alertId, commitAfterSend(state, decision, delivery)); + return; + } + if (decision.kind === "commit") { + await store.save(input.alertId, decision.next); + return; + } + if (state) await store.save(input.alertId, state); + }, +}); diff --git a/apps/api/src/monitoring/notifier.ts b/apps/api/src/monitoring/notifier.ts index 34df23e..0261a9b 100644 --- a/apps/api/src/monitoring/notifier.ts +++ b/apps/api/src/monitoring/notifier.ts @@ -1,5 +1,9 @@ import nodemailer from "nodemailer"; -import { config } from "../utils/config"; +import { getSmtpSettings } from "../db/repo/settings"; +import type { FailureNotificationContext, MailDelivery } from "../types"; +import { appBaseUrl } from "../utils/routes"; +import { buildSlackMessage } from "./slack-message"; +import { type AlertDetails, buildDeploymentFailureEmail, buildEmail } from "./templates"; type NotifyOpts = { channel: string; @@ -8,45 +12,66 @@ type NotifyOpts = { alertType: string; threshold: number | null; currentValue: number; + details?: AlertDetails; +}; + +type SmtpConfig = { + host: string; + port: number; + user: string; + pass: string; + from: string; +}; + +const loadSmtpConfig = async (): Promise => { + try { + const db = await getSmtpSettings(); + if (db?.host) { + return { host: db.host, port: db.port, user: db.user, pass: db.pass, from: db.fromAddress }; + } + } catch {} + return null; }; let transporter: nodemailer.Transporter | null = null; +let transporterKey = ""; -const getTransporter = () => { - if (transporter) return transporter; - if (!config.smtpHost) return null; +const getTransporter = async (smtp: SmtpConfig) => { + const key = `${smtp.host}:${smtp.port}:${smtp.user}`; + if (transporter && transporterKey === key) return transporter; transporter = nodemailer.createTransport({ - host: config.smtpHost, - port: config.smtpPort, - secure: config.smtpPort === 465, - auth: config.smtpUser && config.smtpPass ? { user: config.smtpUser, pass: config.smtpPass } : undefined, + host: smtp.host, + port: smtp.port, + secure: smtp.port === 465, + auth: smtp.user && smtp.pass ? { user: smtp.user, pass: smtp.pass } : undefined, + connectionTimeout: 10_000, + greetingTimeout: 10_000, + socketTimeout: 30_000, }); + transporterKey = key; return transporter; }; -const subject = (projectName: string, alertType: string) => - `[Dequel] ${alertType.toUpperCase()} alert — ${projectName}`; - -const textBody = (projectName: string, alertType: string, threshold: number | null, currentValue: number) => - `Alert: ${projectName}\n\nType: ${alertType}\nThreshold: ${threshold ?? "N/A"}\nCurrent value: ${currentValue}\n\nThis is an automated notification from Dequel.`; - const sendEmail = async ( to: string, projectName: string, alertType: string, threshold: number | null, currentValue: number, + details?: AlertDetails, ) => { - const t = getTransporter(); - if (!t) { + const smtp = await loadSmtpConfig(); + if (!smtp) { console.warn(`[Notifier] SMTP not configured — skipping email to ${to}`); return; } + const { subject, html } = buildEmail(alertType, projectName, threshold, currentValue, details); + const t = await getTransporter(smtp); await t.sendMail({ - from: config.smtpFrom, + from: smtp.from, to, - subject: subject(projectName, alertType), - text: textBody(projectName, alertType, threshold, currentValue), + subject, + html, }); }; @@ -56,27 +81,16 @@ const sendSlack = async ( alertType: string, threshold: number | null, currentValue: number, + details?: AlertDetails, ) => { - const payload = { - text: subject(projectName, alertType), - blocks: [ - { type: "header", text: { type: "plain_text", text: `⚠️ Dequel Alert: ${projectName}` } }, - { - type: "section", - fields: [ - { type: "mrkdwn", text: `*Type:* ${alertType}` }, - { type: "mrkdwn", text: `*Threshold:* ${threshold ?? "N/A"}` }, - { type: "mrkdwn", text: `*Current:* ${currentValue}` }, - ], - }, - ], - }; + const payload = buildSlackMessage(alertType, projectName, threshold, currentValue, details); const res = await fetch(webhookUrl, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), + signal: AbortSignal.timeout(10_000), }); - if (!res.ok) console.warn(`[Notifier] Slack webhook returned ${res.status}`); + if (!res.ok) throw new Error(`Slack webhook returned ${res.status}`); }; const sendWebhook = async ( @@ -98,26 +112,55 @@ const sendWebhook = async ( method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), + signal: AbortSignal.timeout(10_000), }); - if (!res.ok) console.warn(`[Notifier] Webhook returned ${res.status}`); + if (!res.ok) throw new Error(`Webhook returned ${res.status}`); }; -export const sendNotification = async (opts: NotifyOpts): Promise => { - const { channel, destination, projectName, alertType, threshold, currentValue } = opts; +export const sendNotification = async (opts: NotifyOpts): Promise => { + const { channel, projectName, alertType, threshold, currentValue, details } = opts; + let destination = opts.destination; + if (!destination && channel === "email") { + const smtp = await loadSmtpConfig(); + destination = smtp?.from ?? null; + } if (!destination) { console.warn(`[Notifier] No destination for ${channel} alert — skipping`); - return; - } - const fn = - channel === "email" ? sendEmail : channel === "slack" ? sendSlack : channel === "webhook" ? sendWebhook : null; - if (!fn) { - console.warn(`[Notifier] Unknown channel: ${channel}`); - return; + return { status: "skipped", reason: "no_recipient" }; } try { - await fn(destination, projectName, alertType, threshold, currentValue); + if (channel === "email") { + await sendEmail(destination, projectName, alertType, threshold, currentValue, details); + } else if (channel === "slack") { + await sendSlack(destination, projectName, alertType, threshold, currentValue, details); + } else if (channel === "webhook") { + await sendWebhook(destination, projectName, alertType, threshold, currentValue); + } else { + console.warn(`[Notifier] Unknown channel: ${channel}`); + return { status: "failed", error: `unknown channel: ${channel}` }; + } console.log(`[Notifier] ${channel} alert sent to ${destination}`); + return { status: "sent" }; } catch (err) { console.error(`[Notifier] Failed to send ${channel} alert:`, err); + return { status: "failed", error: err instanceof Error ? err.message : String(err) }; + } +}; + +export const isSmtpConfigured = async (): Promise => (await loadSmtpConfig()) !== null; + +export const sendDeploymentFailureEmail = async (ctx: FailureNotificationContext): Promise => { + const smtp = await loadSmtpConfig(); + if (!smtp) return { status: "skipped", reason: "no_smtp" }; + const to = smtp.from; + if (!to) return { status: "skipped", reason: "no_recipient" }; + try { + const logsUrl = ctx.projectId ? `${appBaseUrl()}/project/${ctx.projectId}?tab=deployments` : undefined; + const { subject, html } = buildDeploymentFailureEmail({ ...ctx, logsUrl }); + const t = await getTransporter(smtp); + await t.sendMail({ from: smtp.from, to, subject, html }); + return { status: "sent" }; + } catch (err) { + return { status: "failed", error: err instanceof Error ? err.message : String(err) }; } }; diff --git a/apps/api/src/monitoring/project-status.ts b/apps/api/src/monitoring/project-status.ts new file mode 100644 index 0000000..064e5ff --- /dev/null +++ b/apps/api/src/monitoring/project-status.ts @@ -0,0 +1,87 @@ +import { getProjectById, getServerById, listDeployments, listProjectEvents, listRoutesByDeployment } from "../db/repo"; +import { scalingEngine } from "../scaling/engine"; +import type { ProjectStatus, StatusFailure, StatusHistoryEntry } from "../types"; +import { collectContainerStats } from "./container-stats"; +import { getHttpErrorRate } from "./http-error-rate"; +import { evaluateHealth } from "./health"; + +export const getProjectStatus = async (projectId: string, windowSeconds: number): Promise => { + const project = await getProjectById(projectId); + if (!project) throw new Error(`Project ${projectId} not found`); + + const sinceIso = new Date(Date.now() - windowSeconds * 1000).toISOString(); + const deployments = await listDeployments(projectId); + const running = deployments.filter((d) => d.status === "running"); + + const [events, replicas, containers, http] = await Promise.all([ + listProjectEvents(projectId, sinceIso), + scalingEngine.getProjectReplicas(projectId), + collectContainerStats(deployments), + getHttpErrorRate(projectId, windowSeconds), + ]); + + const serverId = running[0]?.serverId ?? project.serverId ?? null; + const server = serverId && serverId !== "local" ? await getServerById(serverId).catch(() => null) : null; + + const routeErrors: string[] = []; + for (const dep of running) { + const routes = await listRoutesByDeployment(dep.id).catch(() => []); + for (const route of routes) { + if (route.lastError) routeErrors.push(route.lastError); + else if (route.status === "failed") routeErrors.push(`Route ${route.hostname} is ${route.status}`); + } + } + + const failures: StatusFailure[] = events + .filter((e) => e.type === "failed") + .map((e) => ({ + deploymentId: e.deploymentId, + message: e.message, + commitSha: e.commitSha, + sourceRef: e.sourceRef, + finishedAt: e.finishedAt, + recovered: e.deploymentStatus !== "failed", + notifiedAt: e.sentAt, + })); + + const history: StatusHistoryEntry[] = events.map((e) => ({ + deploymentId: e.deploymentId, + type: e.type, + message: e.message, + at: e.createdAt, + })); + + const health = evaluateHealth({ + server: server + ? { + status: server.status, + lastHeartbeatAgeMs: server.lastHeartbeat ? Date.now() - new Date(server.lastHeartbeat).getTime() : null, + } + : null, + runningDeployments: running.length, + hasDeployments: deployments.length > 0, + routeErrors, + http: { available: http.available, errorRate: http.errorRate }, + }); + + return { + projectId, + windowSeconds, + health, + failures, + history, + replicas, + resources: { + server: server + ? { + status: server.status, + cpuUsedPercent: server.cpuUsedPercent, + memoryTotalMb: server.memoryTotalMb, + lastHeartbeatAt: server.lastHeartbeat, + } + : null, + containers, + }, + http, + }; +}; diff --git a/apps/api/src/monitoring/slack-message.ts b/apps/api/src/monitoring/slack-message.ts new file mode 100644 index 0000000..7619b08 --- /dev/null +++ b/apps/api/src/monitoring/slack-message.ts @@ -0,0 +1,75 @@ +import type { AlertDetails } from "./templates"; + +const formatValue = (alertType: string, value: number): string => { + if (alertType === "downtime") return "Service down"; + if (alertType === "memory") return `${value.toFixed(0)} MB`; + return `${value.toFixed(1)}%`; +}; + +const formatThreshold = (alertType: string, threshold: number | null): string => { + if (alertType === "downtime" || threshold === null) return "N/A"; + return alertType === "memory" ? String(threshold) : `${threshold}%`; +}; + +const formatContainer = (alertType: string, c: { name: string; value: number }): string => + `\`${c.name}\` ${alertType === "memory" ? `${c.value.toFixed(0)} MB` : `${c.value.toFixed(1)}%`}`; + +type SlackBlock = + | { type: "header"; text: { type: string; text: string } } + | { + type: "section"; + text?: { type: string; text: string }; + fields?: { type: string; text: string }[]; + } + | { type: "actions"; elements: { type: string; text: { type: string; text: string }; url: string }[] }; + +export const buildSlackMessage = ( + alertType: string, + projectName: string, + threshold: number | null, + currentValue: number, + details?: AlertDetails, +): { text: string; blocks: SlackBlock[] } => { + const blocks: SlackBlock[] = [ + { type: "header", text: { type: "plain_text", text: `Dequel alert: ${projectName}` } }, + { + type: "section", + fields: [ + { type: "mrkdwn", text: `*Type:* ${alertType.replace("_", " ")}` }, + { type: "mrkdwn", text: `*Threshold:* ${formatThreshold(alertType, threshold)}` }, + { type: "mrkdwn", text: `*Current:* ${formatValue(alertType, currentValue)}` }, + ], + }, + ]; + const containers = details?.containers ?? []; + if (containers.length > 0) { + blocks.push({ + type: "section", + text: { + type: "mrkdwn", + text: containers.map((c) => `• ${formatContainer(alertType, c)}`).join("\n"), + }, + }); + } + if (details?.scaling) { + const s = details.scaling; + const title = s.kind === "enable_autoscaling" ? "Autoscaling is off" : "Replica limit reached"; + const body = + s.kind === "enable_autoscaling" + ? "Enable autoscaling so Dequel adds capacity automatically while load stays high." + : `At the replica limit (${s.current}/${s.maxReplicas}) — raise max replicas so autoscaling can add capacity.`; + const cta = s.kind === "enable_autoscaling" ? "Set up autoscaling" : "Adjust scaling"; + blocks.push({ type: "section", text: { type: "mrkdwn", text: `*${title}*\n${body}` } }); + blocks.push({ + type: "actions", + elements: [{ type: "button", text: { type: "plain_text", text: cta }, url: s.url }], + }); + } + if (details?.projectUrl) { + blocks.push({ + type: "actions", + elements: [{ type: "button", text: { type: "plain_text", text: "Open project" }, url: details.projectUrl }], + }); + } + return { text: `${projectName}: ${alertType} alert`, blocks }; +}; diff --git a/apps/api/src/monitoring/templates.ts b/apps/api/src/monitoring/templates.ts new file mode 100644 index 0000000..e8c8868 --- /dev/null +++ b/apps/api/src/monitoring/templates.ts @@ -0,0 +1,216 @@ +import { readFileSync } from "node:fs"; +import { join } from "node:path"; +import type { ScalingSuggestion } from "./alert-guard"; + +export interface AlertDetails { + containers?: { name: string; value: number }[]; + lastRunningAt?: string | null; + logsUrl?: string; + appUrl?: string; + projectUrl?: string; + scaling?: (ScalingSuggestion & { url: string }) | null; +} + +export const EMAIL_ATTACHMENTS = [ + { + filename: "logo.webp", + path: join(import.meta.dir, "templates", "logo.webp"), + cid: "dequel-logo", + }, +]; + +const THEME = { + cardBorder: "#e4e4e7", + text: "#18181b", + textMuted: "#71717a", + accent: "#ea580c", + link: "#7c3aed", + red: "#dc2626", + redBg: "#fef2f2", + amber: "#d97706", + amberBg: "#fff7ed", + green: "#059669", + greenBg: "#ecfdf5", +}; + +const loadHtml = (filename: string): string => { + const filepath = join(import.meta.dir, "templates", filename); + return readFileSync(filepath, "utf-8"); +}; + +const layoutHtml = loadHtml("layout.html"); +const deployFailureTpl = loadHtml("deploy-failure.html"); +const alertMetricTpl = loadHtml("alert-metric.html"); +const alertDowntimeTpl = loadHtml("alert-downtime.html"); +const alertCertExpiryTpl = loadHtml("alert-cert-expiry.html"); +const alertDefaultTpl = loadHtml("alert-default.html"); +const smtpTestTpl = loadHtml("smtp-test.html"); + +const truncated = (s: string, n = 160) => (s.length > n ? `${s.slice(0, n)}…` : s); + +const renderBadge = (text: string, bg: string, textColor: string) => { + if (!text) return ""; + return `
+ ${text} +
`; +}; + +const renderButton = (href?: string, text?: string) => { + if (!href || !text) return ""; + return `
+ ${text} +
`; +}; + +const renderRow = (label: string, value: string) => ` +
+ ${label} + ${value} +
`; + +const renderScalingSuggestion = (s: AlertDetails["scaling"]): string => { + if (!s) return ""; + const enable = s.kind === "enable_autoscaling"; + const title = enable ? "Autoscaling is off" : "Replica limit reached"; + const text = enable + ? "Enable autoscaling for this project and Dequel will add capacity automatically while load stays high." + : `This project is at its replica limit (${s.current}/${s.maxReplicas}). Raise max replicas so autoscaling can add capacity.`; + const cta = enable ? "Set up autoscaling" : "Adjust scaling"; + return `
+
${title}
+
${text}
+ ${cta} +
`; +}; + +const renderShell = (badgeText: string, badgeBg: string, badgeTextColor: string, bodyContent: string): string => { + const badgeHtml = renderBadge(badgeText, badgeBg, badgeTextColor); + return layoutHtml.replace("{{BADGE_HTML}}", badgeHtml).replace("{{BODY_CONTENT}}", bodyContent); +}; + +export const buildEmail = ( + alertType: string, + projectName: string, + threshold: number | null, + currentValue: number, + details?: AlertDetails, +): { subject: string; html: string } => { + switch (alertType) { + case "cpu": + case "memory": { + const isCpu = alertType === "cpu"; + const unit = "%"; + const containers = details?.containers ?? []; + const containerRows = containers.map((c) => renderRow(c.name, `${c.value.toFixed(1)}${unit}`)).join(""); + const actionButton = details?.appUrl ? renderButton(details.appUrl, "View Application") : ""; + + const body = alertMetricTpl + .replace(/{{METRIC_TYPE}}/g, isCpu ? "CPU" : "Memory") + .replace(/{{PROJECT_NAME}}/g, projectName) + .replace(/{{CURRENT_VALUE}}/g, currentValue.toFixed(1)) + .replace(/{{THRESHOLD}}/g, String(threshold ?? "N/A")) + .replace(/{{UNIT}}/g, unit) + .replace("{{CONTAINER_ROWS}}", containerRows) + .replace("{{SCALING_HTML}}", renderScalingSuggestion(details?.scaling)) + .replace("{{ACTION_BUTTON}}", actionButton); + + return { + subject: `${isCpu ? "High CPU" : "High memory"} on ${projectName} (${currentValue.toFixed(0)}${unit})`, + html: renderShell(isCpu ? "CPU ALERT" : "MEMORY ALERT", THEME.amberBg, THEME.amber, body), + }; + } + + case "downtime": { + const lastRunningRow = details?.lastRunningAt + ? renderRow("Last running", new Date(details.lastRunningAt).toUTCString()) + : ""; + const actionButton = details?.logsUrl ? renderButton(details.logsUrl, "View Logs") : ""; + + const body = alertDowntimeTpl + .replace(/{{PROJECT_NAME}}/g, projectName) + .replace("{{LAST_RUNNING_ROW}}", lastRunningRow) + .replace("{{ACTION_BUTTON}}", actionButton); + + return { + subject: `${projectName} is down`, + html: renderShell("SERVICE DOWN", THEME.redBg, THEME.red, body), + }; + } + + case "cert_expiry": { + const daysRow = renderRow("Days remaining", String(currentValue)); + const actionButton = details?.appUrl ? renderButton(details.appUrl, "View Site") : ""; + + const body = alertCertExpiryTpl + .replace(/{{PROJECT_NAME}}/g, projectName) + .replace(/{{CURRENT_VALUE}}/g, String(currentValue)) + .replace("{{DAYS_ROW}}", daysRow) + .replace("{{ACTION_BUTTON}}", actionButton); + + return { + subject: `SSL certificate for ${projectName} expires in ${currentValue} days`, + html: renderShell("CERTIFICATE EXPIRY", THEME.amberBg, THEME.amber, body), + }; + } + + default: { + const detailsRows = [ + renderRow("Type", alertType), + renderRow("Threshold", threshold !== null ? String(threshold) : "N/A"), + renderRow("Current value", String(currentValue)), + ].join(""); + + const body = alertDefaultTpl.replace(/{{PROJECT_NAME}}/g, projectName).replace("{{DETAILS_ROWS}}", detailsRows); + + return { + subject: `[Dequel] ${alertType} alert — ${projectName}`, + html: renderShell("ALERT", THEME.amberBg, THEME.amber, body), + }; + } + } +}; + +export const buildDeploymentFailureEmail = (ctx: { + projectName: string; + failureReason: string | null; + commitSha: string | null; + sourceRef: string; + finishedAt: string | null; + logsUrl?: string; +}): { subject: string; html: string } => { + const commitItem = ctx.commitSha + ? `
  • Commit: ${ + ctx.failureReason + ? `${truncated(ctx.failureReason, 140)} (${ctx.commitSha.slice(0, 12)})` + : `${ctx.commitSha.slice(0, 12)}` + }
  • ` + : ""; + const sourceItem = ctx.sourceRef + ? `
  • Source: ${ctx.sourceRef}
  • ` + : ""; + const reasonItem = + ctx.failureReason && !ctx.commitSha + ? `
  • Reason: ${truncated(ctx.failureReason, 160)}
  • ` + : ""; + + const detailsItems = `${commitItem}${sourceItem}${reasonItem}`; + const actionButton = ctx.logsUrl ? renderButton(ctx.logsUrl, "View Logs") : ""; + + const body = deployFailureTpl + .replace(/{{PROJECT_NAME}}/g, ctx.projectName) + .replace(/{{LOGS_URL}}/g, ctx.logsUrl || "#") + .replace("{{DETAILS_ITEMS}}", detailsItems) + .replace("{{ACTION_BUTTON}}", actionButton); + + return { + subject: `deploy failed for ${ctx.projectName}`, + html: renderShell("DEPLOY FAILED", THEME.redBg, THEME.red, body), + }; +}; + +export const buildSmtpTestEmail = (): { subject: string; html: string } => { + return { + subject: "[Dequel] SMTP Test Email", + html: renderShell("SMTP TEST", THEME.greenBg, THEME.green, smtpTestTpl), + }; +}; diff --git a/apps/api/src/monitoring/templates/alert-cert-expiry.html b/apps/api/src/monitoring/templates/alert-cert-expiry.html new file mode 100644 index 0000000..9090505 --- /dev/null +++ b/apps/api/src/monitoring/templates/alert-cert-expiry.html @@ -0,0 +1,8 @@ +

    + The SSL certificate for {{PROJECT_NAME}} expires in {{DAYS_REMAINING}} days. +

    +{{DAYS_ROW}} +{{ACTION_BUTTON}} +

    + Learn more about SSL certificates on Dequel. +

    diff --git a/apps/api/src/monitoring/templates/alert-default.html b/apps/api/src/monitoring/templates/alert-default.html new file mode 100644 index 0000000..ff916ca --- /dev/null +++ b/apps/api/src/monitoring/templates/alert-default.html @@ -0,0 +1,7 @@ +

    + An alert was triggered for {{PROJECT_NAME}}. +

    +{{DETAILS_ROWS}} +

    + Learn more about alerts on Dequel. +

    diff --git a/apps/api/src/monitoring/templates/alert-downtime.html b/apps/api/src/monitoring/templates/alert-downtime.html new file mode 100644 index 0000000..83a82ce --- /dev/null +++ b/apps/api/src/monitoring/templates/alert-downtime.html @@ -0,0 +1,12 @@ +

    + {{PROJECT_NAME}} has no running deployments. The service appears to be down. +

    +
    +
    Offline
    +
    0 running deployments
    +
    +{{LAST_RUNNING_ROW}} +{{ACTION_BUTTON}} +

    + Learn more about service availability on Dequel. +

    diff --git a/apps/api/src/monitoring/templates/alert-metric.html b/apps/api/src/monitoring/templates/alert-metric.html new file mode 100644 index 0000000..118a2a5 --- /dev/null +++ b/apps/api/src/monitoring/templates/alert-metric.html @@ -0,0 +1,13 @@ +

    + {{METRIC_TYPE}} usage on {{PROJECT_NAME}} has exceeded your configured threshold. +

    +
    +
    {{CURRENT_VALUE}}{{UNIT}}
    +
    Current {{METRIC_TYPE}} — threshold {{THRESHOLD}}{{UNIT}}
    +
    +{{CONTAINER_ROWS}} +{{SCALING_HTML}} +{{ACTION_BUTTON}} +

    + Learn more about monitoring on Dequel. +

    diff --git a/apps/api/src/monitoring/templates/deploy-failure.html b/apps/api/src/monitoring/templates/deploy-failure.html new file mode 100644 index 0000000..cfcb7dc --- /dev/null +++ b/apps/api/src/monitoring/templates/deploy-failure.html @@ -0,0 +1,11 @@ +

    + We encountered an error during the deploy process for {{PROJECT_NAME}}. + This means your deploy didn't complete successfully and your latest changes may not be live. +

    +
      + {{DETAILS_ITEMS}} +
    +{{ACTION_BUTTON}} +

    + Learn more about troubleshooting deploys on Dequel. +

    diff --git a/apps/api/src/monitoring/templates/layout.html b/apps/api/src/monitoring/templates/layout.html new file mode 100644 index 0000000..2fb0743 --- /dev/null +++ b/apps/api/src/monitoring/templates/layout.html @@ -0,0 +1,116 @@ + + + + + + + Dequel Notification + + + + + + + +
    + + + + + + + + + + + + + + + +
    + + + + + +
    + + + + + +
    + Dequel Logo + + Dequel +
    +
    + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
               
               
               
    +
    +
    + {{BADGE_HTML}} + {{BODY_CONTENT}} +
    +

    + The Dequel team +

    +

    + Don't want to receive these emails? You can change your notification settings for your workspace or just this service. +

    +

    + Need more help? Contact our Support team. +

    +
    + © 2026 Dequel +
    +
    +
    + + + \ No newline at end of file diff --git a/apps/api/src/monitoring/templates/logo.webp b/apps/api/src/monitoring/templates/logo.webp new file mode 100644 index 0000000..1df87a6 Binary files /dev/null and b/apps/api/src/monitoring/templates/logo.webp differ diff --git a/apps/api/src/monitoring/templates/smtp-test.html b/apps/api/src/monitoring/templates/smtp-test.html new file mode 100644 index 0000000..197ada5 --- /dev/null +++ b/apps/api/src/monitoring/templates/smtp-test.html @@ -0,0 +1,11 @@ +

    + This is a test email sent from your Dequel instance. + If you are receiving this, your SMTP settings have been configured successfully. +

    +
    +
    SMTP Transport Verified
    +
    Notifications and deployment alerts are ready to be delivered via email.
    +
    +

    + Learn more about managing notifications on Dequel. +

    diff --git a/apps/api/src/orchestrator/__tests__/pipeline-cleanup.test.ts b/apps/api/src/orchestrator/__tests__/pipeline-cleanup.test.ts index e9330d7..311c3ee 100644 --- a/apps/api/src/orchestrator/__tests__/pipeline-cleanup.test.ts +++ b/apps/api/src/orchestrator/__tests__/pipeline-cleanup.test.ts @@ -48,6 +48,8 @@ const mockDb = { listEnvironmentVariablesForDeploy: mock(() => Promise.resolve([])), listVolumes: mock(() => Promise.resolve([])), listDeployments: mock(() => Promise.resolve([])), + listProjectEvents: mock(() => Promise.resolve([])), + listRoutesByDeployment: mock(() => Promise.resolve([])), listAllDatabases: mock(() => Promise.resolve([])), deleteDeploymentAndLogs: mock(() => Promise.resolve()), getScalingPolicy: mock(() => Promise.resolve(null)), @@ -65,11 +67,14 @@ const mockDb = { updateDomainValidation: mock(() => Promise.resolve()), listDomains: mock(() => Promise.resolve([])), createDeploymentEvent: mock(() => Promise.resolve()), + recordDeploymentFailure: mock(() => Promise.resolve({ outcome: "recorded" })), + recordDeploymentCancellation: mock(() => Promise.resolve({ outcome: "recorded" })), }; mock.module(fileUrl("../../db/repo"), () => mockDb); mock.module(fileUrl("../runtime"), () => ({ + run: mock(() => Promise.resolve("")), deployContainer: mock(() => Promise.resolve({ containerName: "test-project-abc12345", diff --git a/apps/api/src/orchestrator/pipeline.ts b/apps/api/src/orchestrator/pipeline.ts index 7ebd033..7f0dc5b 100644 --- a/apps/api/src/orchestrator/pipeline.ts +++ b/apps/api/src/orchestrator/pipeline.ts @@ -13,10 +13,14 @@ import { listDeployments, listEnvironmentVariablesForDeploy, listVolumes, + recordDeploymentCancellation, + recordDeploymentFailure, updateDeploymentCommitSha, updateDeploymentStatus, } from "../db/repo"; import { deployments } from "../db/schema"; +import { config } from "../utils/config"; +import { CANCELLED_FAILURE_REASON } from "../utils/failure-outcome"; import { ensureProjectDashboard } from "../utils/grafana"; import { buildWithCompose, destroyComposeStack } from "./compose"; import { deployComposeStack } from "./compose-deploy"; @@ -88,13 +92,12 @@ export class PipelineOrchestrator { await Promise.all([ this.queue.remove(deploymentId), - updateDeploymentStatus(deploymentId, "failed", { failureReason: "Cancelled" }), - appendLog(deploymentId, "system", "Deployment cancelled by user"), - createDeploymentEvent({ + recordDeploymentCancellation({ deploymentId, - type: "cancelled", - message: "Deployment cancelled by user", + reason: CANCELLED_FAILURE_REASON, + source: "pipeline", }), + appendLog(deploymentId, "system", "Deployment cancelled by user"), ]); logBus.publish({ deploymentId, @@ -466,13 +469,7 @@ export class PipelineOrchestrator { const message = summarizeDeploymentError(error); console.error(`[Orchestrator] Deployment ${deploymentId} failed:`, error); await emitLog(deploymentId, "system", `Deployment failed: ${message}`); - await updateDeploymentStatus(deploymentId, "failed", { failureReason: message }); - await createDeploymentEvent({ - deploymentId, - type: "failed", - message, - metadata: { stage: "unknown" }, - }); + await recordDeploymentFailure({ deploymentId, reason: message, source: "pipeline" }); if (!deployed) { await emitLog(deploymentId, "system", "Cleaning up Docker resources from failed deployment"); @@ -604,7 +601,7 @@ export class PipelineOrchestrator { const message = summarizeDeploymentError(error); console.error(`[Orchestrator] Rollback of ${targetDeploymentId} failed:`, error); await emitLog(targetDeploymentId, "system", `Rollback failed: ${message}`); - await updateDeploymentStatus(targetDeploymentId, "failed", { failureReason: message }); + await recordDeploymentFailure({ deploymentId: targetDeploymentId, reason: message, source: "rollback" }); throw error; } } diff --git a/apps/api/src/orchestrator/reconciliation.ts b/apps/api/src/orchestrator/reconciliation.ts index 278c59e..1261f71 100644 --- a/apps/api/src/orchestrator/reconciliation.ts +++ b/apps/api/src/orchestrator/reconciliation.ts @@ -1,6 +1,7 @@ import { and, eq, lt } from "drizzle-orm"; import { getDb } from "../db/db-provider"; -import { agentJobs, deployments, servers } from "../db/schema"; +import { recordDeploymentFailure } from "../db/repo"; +import { agentJobs, servers } from "../db/schema"; const LEASE_RECOVERY_INTERVAL_MS = 30_000; const STALE_AGENT_THRESHOLD_MS = 5 * 60 * 1000; @@ -113,15 +114,11 @@ const cleanAbandonedJobs = async () => { .execute(); if (job.deploymentId) { - await db - .update(deployments) - .set({ - status: "failed", - failureReason: "Agent job abandoned", - finishedAt: new Date(), - }) - .where(eq(deployments.id, job.deploymentId)) - .execute(); + await recordDeploymentFailure({ + deploymentId: job.deploymentId, + reason: "Agent job abandoned", + source: "reconciler", + }).catch((err) => console.error(`[Reconciliation] Failed to record failure for ${job.deploymentId}:`, err)); } console.log(`[Reconciliation] Cleaned abandoned job ${job.id} (started ${job.startedAt})`); diff --git a/apps/api/src/scaling/engine.ts b/apps/api/src/scaling/engine.ts index c4a7b48..a4faafc 100644 --- a/apps/api/src/scaling/engine.ts +++ b/apps/api/src/scaling/engine.ts @@ -6,6 +6,7 @@ import type { Server } from "../types"; import { config } from "../utils/config"; import { DEQUEL_MANAGED_LABEL } from "../utils/dequel-labels"; import { dockerBin } from "../utils/docker-bin"; +import { slugify } from "../utils/routes"; import { execDockerSshCommand, syncRemoteCaddyRoute } from "../utils/ssh"; import { run, tryRun } from "./docker-utils"; @@ -232,7 +233,10 @@ class ScalingEngine { const containers = await agentStatsCache.get(target.server!.id); let count = 0; for (const stat of containers.values()) { - if (stat.replica && dep.id && stat.deploymentId === dep.id) count++; + if (!stat.replica || !dep.id) continue; + if (stat.containerName === dep.containerName || stat.containerName.startsWith(`deploy-${dep.id}`)) { + count++; + } } if (containers.has(dep.containerName ?? "")) count++; return Math.max(1, count); @@ -242,7 +246,7 @@ class ScalingEngine { "ps", "-q", "--filter", - "label=com.dequel.managed=1", + `label=${DEQUEL_MANAGED_LABEL}`, "--filter", `name=deploy-${dep.id}-replica-`, ]); @@ -269,7 +273,10 @@ class ScalingEngine { const containers = new Set(); for (const m of matches) { const parts = m.replace("reverse_proxy", "").trim().split(/\s+/); - for (const p of parts) containers.add(p.split(":")[0]); + for (const p of parts) { + if (!p.includes(":") || p.startsWith("{")) continue; + containers.add(p.split(":")[0]); + } } return Math.max(1, containers.size); } catch { @@ -277,6 +284,21 @@ class ScalingEngine { } } + async getProjectReplicas(projectId: string): Promise<{ current: number } | null> { + try { + const deployments = await listDeployments(projectId); + const runningDep = deployments.find((d) => d.status === "running") ?? deployments.find((d) => d.containerName); + if (!runningDep) return null; + const project = await getProjectById(projectId); + const target = await this.resolveTarget(runningDep); + const slug = project ? slugify(project.name) : projectId; + return { current: await this.getCurrentReplicas(slug, target, runningDep) }; + } catch (err) { + console.warn(`[Scaling] Failed to get replicas for project ${projectId}:`, err); + return null; + } + } + private async scaleUp( dep: { id: string; projectId: string | null; containerName: string | null; serverId?: string | null }, maxReplicas: number, @@ -499,6 +521,7 @@ class ScalingEngine { await import("../utils/ingress"); const ingressServer = await getIngressServer(); const viaIngress = shouldRouteViaIngress(target.server ?? null, ingressServer); + const { caddyReverseProxy, caddySite } = await import("../utils/caddy-site"); const caddySnippet = viaIngress ? projectServerSite( `${slug}.${baseDomain}`, @@ -506,7 +529,7 @@ class ScalingEngine { targets.map((t) => t.split(":")[0]), true, ) - : `${slug}.${baseDomain} {\n reverse_proxy ${targets.join(" ")} {\n header_up Host {upstream_hostport}\n }\n}\n`; + : caddySite(`${slug}.${baseDomain}`, caddyReverseProxy(targets.join(" "))); if (target.mode === "ssh") { await syncRemoteCaddyRoute(target.server!, `${slug}.caddy`, caddySnippet); diff --git a/apps/api/src/types.ts b/apps/api/src/types.ts index 4f1ef2b..655e8d6 100644 --- a/apps/api/src/types.ts +++ b/apps/api/src/types.ts @@ -16,7 +16,7 @@ export type SslStatus = "pending" | "provisioned" | "failed"; export type ServerStatus = "pending" | "connected" | "disconnected" | "failed"; export type ServerMode = "local" | "ssh" | "agent" | "docker_tcp"; export type AlertChannel = "email" | "slack" | "webhook"; -export type AlertType = "cpu" | "memory" | "error_rate" | "downtime" | "cert_expiry"; +export type AlertType = "cpu" | "memory" | "downtime" | "cert_expiry"; export interface Project { id: string; @@ -290,6 +290,92 @@ export interface CreateAlertInput { destination?: string; } +export type FailureSource = "pipeline" | "rollback" | "ssh" | "agent" | "job-channel" | "reconciler" | "dispatch"; + +export interface RecordFailureInput { + deploymentId: string; + reason: string; + source: FailureSource; + cancel?: boolean; +} + +export interface RecordFailureOutcome { + claimed: boolean; + eventId: string | null; +} + +export interface FailureNotificationContext { + eventId: string; + deploymentId: string; + projectId: string | null; + projectName: string; + failureReason: string | null; + commitSha: string | null; + sourceRef: string; + finishedAt: string | null; + attempt: number; +} + +export type MailDelivery = + | { status: "sent" } + | { status: "skipped"; reason: "no_smtp" | "no_recipient" } + | { status: "failed"; error: string }; + +export type HealthStatus = "ok" | "warn" | "fail" | "unknown"; +export type OverallHealth = "healthy" | "degraded" | "down"; + +export interface HealthCheck { + name: "server" | "containers" | "ingress" | "http"; + status: HealthStatus; + detail: string | null; +} + +export interface StatusFailure { + deploymentId: string; + message: string | null; + commitSha: string | null; + sourceRef: string; + finishedAt: string | null; + recovered: boolean; + notifiedAt: string | null; +} + +export interface StatusHistoryEntry { + deploymentId: string; + type: string; + message: string | null; + at: string; +} + +export interface HttpErrorRate { + source: "loki"; + available: boolean; + windowSeconds: number; + totalRequests: number; + errorRequests: number; + errorRate: number | null; + byStatus: { status: string; count: number }[]; +} + +export interface ProjectStatus { + projectId: string; + windowSeconds: number; + health: { overall: OverallHealth; checks: HealthCheck[] }; + failures: StatusFailure[]; + history: StatusHistoryEntry[]; + replicas: { current: number } | null; + resources: { + server: { + status: string; + cpuUsedPercent: number | null; + memoryTotalMb: number | null; + lastHeartbeatAt: string | null; + } | null; + containers: { name: string; cpuPercent: number; memoryMb: number }[]; + }; + http: HttpErrorRate; +} + export interface Deployment { id: string; projectId: string | null; diff --git a/apps/api/src/utils/caddy-site.ts b/apps/api/src/utils/caddy-site.ts new file mode 100644 index 0000000..6898bc7 --- /dev/null +++ b/apps/api/src/utils/caddy-site.ts @@ -0,0 +1,10 @@ +export const CADDY_ACCESS_LOG_BLOCK = ` log { + output stdout + format json + }`; + +export const caddySite = (hosts: string, reverseProxy: string): string => + `${hosts} {\n${CADDY_ACCESS_LOG_BLOCK}\n${reverseProxy}\n}\n`; + +export const caddyReverseProxy = (targets: string): string => + ` reverse_proxy ${targets} {\n header_up Host {upstream_hostport}\n }`; diff --git a/apps/api/src/utils/config-loader.ts b/apps/api/src/utils/config-loader.ts index 58ac104..02a81a6 100644 --- a/apps/api/src/utils/config-loader.ts +++ b/apps/api/src/utils/config-loader.ts @@ -18,11 +18,6 @@ export interface FileConfig { queueConcurrency?: number; queueRetryMax?: number; queueRetryBaseMs?: number; - smtpHost?: string; - smtpPort?: number; - smtpUser?: string; - smtpPass?: string; - smtpFrom?: string; alertEvalIntervalMs?: number; githubClientId?: string; githubClientSecret?: string; diff --git a/apps/api/src/utils/config.ts b/apps/api/src/utils/config.ts index 2acb8a2..495caa6 100644 --- a/apps/api/src/utils/config.ts +++ b/apps/api/src/utils/config.ts @@ -32,12 +32,8 @@ export const config = { queueConcurrency: withFile("QUEUE_CONCURRENCY", "3", Number), queueRetryMax: withFile("QUEUE_RETRY_MAX", "5", Number), queueRetryBaseMs: withFile("QUEUE_RETRY_BASE_MS", "5000", Number), - smtpHost: withFile("SMTP_HOST", ""), - smtpPort: withFile("SMTP_PORT", "587", Number), - smtpUser: withFile("SMTP_USER", ""), - smtpPass: withFile("SMTP_PASS", ""), - smtpFrom: withFile("SMTP_FROM", "dequel@localhost"), alertEvalIntervalMs: withFile("ALERT_EVAL_INTERVAL_MS", "60000", Number), + failureSweepIntervalMs: withFile("FAILURE_SWEEP_INTERVAL_MS", "60000", Number), githubClientId: withFile("GITHUB_CLIENT_ID", ""), githubClientSecret: withFile("GITHUB_CLIENT_SECRET", ""), githubAppName: withFile("GITHUB_APP_NAME", "Dequel"), diff --git a/apps/api/src/utils/domain-verifier.ts b/apps/api/src/utils/domain-verifier.ts index dc5a08a..05c7d2f 100644 --- a/apps/api/src/utils/domain-verifier.ts +++ b/apps/api/src/utils/domain-verifier.ts @@ -5,6 +5,7 @@ import { getDb } from "../db/db-provider"; import { getProjectById, listDomains, listEnvironmentVariablesForDeploy, updateDomainValidation } from "../db/repo"; import { domains } from "../db/schema"; import { reloadCaddy } from "../orchestrator/runtime"; +import { caddyReverseProxy, caddySite } from "./caddy-site"; import { config } from "./config"; import { resolveServerIp, validateDomain } from "./dns"; @@ -205,9 +206,7 @@ export const buildCaddySnippet = async ( } } const tPort = d.targetPort || port; - customBlocks.push( - `${entryDomain} {\n log {\n output stdout\n format json\n }\n reverse_proxy ${targetContainer}:${tPort} {\n header_up Host {upstream_hostport}\n }\n}\n`, - ); + customBlocks.push(caddySite(entryDomain, caddyReverseProxy(`${targetContainer}:${tPort}`))); } else { if (!defaultDomains.includes(entryDomain)) defaultDomains.push(entryDomain); } @@ -246,7 +245,7 @@ export const buildCaddySnippet = async ( } } - const primaryBlock = `${defaultDomains.join(", ")} {\n log {\n output stdout\n format json\n }\n reverse_proxy ${containerName}:${port} {\n header_up Host {upstream_hostport}\n }\n}\n`; + const primaryBlock = caddySite(defaultDomains.join(", "), caddyReverseProxy(`${containerName}:${port}`)); return [primaryBlock, ...customBlocks].join("\n"); }; diff --git a/apps/api/src/utils/failure-outcome.ts b/apps/api/src/utils/failure-outcome.ts new file mode 100644 index 0000000..9ea15f2 --- /dev/null +++ b/apps/api/src/utils/failure-outcome.ts @@ -0,0 +1,6 @@ +export const CANCELLED_FAILURE_REASON = "Cancelled"; + +export type FailureOutcome = "failure" | "cancelled"; + +export const classifyFailureOutcome = (reason: string | null | undefined): FailureOutcome => + reason === CANCELLED_FAILURE_REASON ? "cancelled" : "failure"; diff --git a/apps/api/src/utils/grafana.ts b/apps/api/src/utils/grafana.ts index b954b0b..2e9fd86 100644 --- a/apps/api/src/utils/grafana.ts +++ b/apps/api/src/utils/grafana.ts @@ -1,5 +1,6 @@ -import { listDomains } from "../db/repo"; import { config } from "./config"; +import { buildProjectRequestHostRegex } from "./loki"; +import { slugify } from "./routes"; interface GrafanaDashboard { dashboard: { @@ -64,24 +65,8 @@ export async function ensureProjectDashboard( projectName: string, containerRegex: string, ): Promise { - const slug = projectName - .toLowerCase() - .replace(/[^a-z0-9-]+/g, "-") - .replace(/^-+|-+$/g, "") - .slice(0, 63); - - const domains = [`${slug}.${config.caddyBaseDomain}`]; - try { - const projectDomains = await listDomains(projectId); - const verified = projectDomains.filter((d) => d.validationStatus === "verified"); - for (const d of verified) { - domains.push(d.domain); - } - } catch (e) { - console.warn("[Grafana] Failed to list domains for dashboard query:", e); - } - - const regexEscaped = domains.map((d) => d.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\\\$&")).join("|"); + const slug = slugify(projectName); + const regexEscaped = await buildProjectRequestHostRegex(projectId); const dashboard: GrafanaDashboard = { dashboard: { diff --git a/apps/api/src/utils/ingress.ts b/apps/api/src/utils/ingress.ts index d969b92..929b04a 100644 --- a/apps/api/src/utils/ingress.ts +++ b/apps/api/src/utils/ingress.ts @@ -1,6 +1,7 @@ import { rm, writeFile } from "node:fs/promises"; import { join } from "node:path"; import { createAgentJob, getPlatformSettings, getServerById, upsertRoute } from "../db/repo"; +import { caddyReverseProxy, caddySite } from "./caddy-site"; import { config } from "./config"; import { removeRemoteCaddyRoute, syncRemoteCaddyRoute } from "./ssh"; @@ -38,10 +39,11 @@ export const projectServerSite = ( viaIngress: boolean, ): string => { const targets = containers.map((c) => `${c}:${port}`).join(" "); + const proxy = caddyReverseProxy(targets); if (viaIngress) { - return `:80 {\n reverse_proxy ${targets} {\n header_up Host {upstream_hostport}\n }\n}\n`; + return caddySite(":80", proxy); } - return `${hostname} {\n reverse_proxy ${targets} {\n header_up Host {upstream_hostport}\n }\n}\n`; + return caddySite(hostname, proxy); }; export const ingressSite = (hostname: string, upstreamHost: string): string => diff --git a/apps/api/src/utils/loki.ts b/apps/api/src/utils/loki.ts new file mode 100644 index 0000000..253835f --- /dev/null +++ b/apps/api/src/utils/loki.ts @@ -0,0 +1,42 @@ +import { getProjectById, listDomains } from "../db/repo"; +import { config } from "./config"; +import { slugify } from "./routes"; + +const LOKI_URL = "http://loki:3100"; +const CADDY_LOG_STREAM = '{container="dequel-caddy-1"}'; + +export const buildProjectRequestHostRegex = async (projectId: string): Promise => { + const project = await getProjectById(projectId); + if (!project) throw new Error(`Project ${projectId} not found`); + const slug = slugify(project.name); + const domains = [`${slug}.${config.caddyBaseDomain}`]; + try { + const projectDomains = await listDomains(projectId); + for (const d of projectDomains) { + if (d.validationStatus === "verified") domains.push(d.domain); + } + } catch (err) { + console.warn("[Loki] Failed to list domains for host regex:", err); + } + return domains.map((d) => d.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\\\$&")).join("|"); +}; + +export const caddyRequestLogSelector = (hostRegex: string): string => + `${CADDY_LOG_STREAM} | json | request_host =~ "^(${hostRegex})$"`; + +export const lokiInstantQuery = async (query: string, timeoutMs = 5000): Promise => { + const url = `${LOKI_URL}/loki/api/v1/query?query=${encodeURIComponent(query)}`; + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), timeoutMs); + try { + const response = await fetch(url, { signal: controller.signal }); + if (!response.ok) return null; + const data = (await response.json()) as any; + if (data.status !== "success") return null; + return data.data?.result ?? null; + } catch { + return null; + } finally { + clearTimeout(timeout); + } +}; diff --git a/apps/api/src/utils/routes.ts b/apps/api/src/utils/routes.ts index 8d9573e..b515cdb 100644 --- a/apps/api/src/utils/routes.ts +++ b/apps/api/src/utils/routes.ts @@ -9,6 +9,9 @@ export const slugify = (s: string) => export const baseDomainFor = () => (config.caddyBaseDomain === "localhost" ? "localhost:80" : config.caddyBaseDomain); +export const appBaseUrl = () => + config.caddyBaseDomain === "localhost" ? "http://localhost" : `https://${config.caddyBaseDomain}`; + export const routeNamesFor = (projectName: string | null, projectId: string | null, deploymentId: string) => { const slug = slugify(projectName || projectId || deploymentId); return { diff --git a/apps/docs/src/content/docs/system-config.md b/apps/docs/src/content/docs/system-config.md index 5851c4f..e258942 100644 --- a/apps/docs/src/content/docs/system-config.md +++ b/apps/docs/src/content/docs/system-config.md @@ -50,29 +50,7 @@ On boot, these values seed the `github_integrations` table. You can also update ### SMTP -Set these to enable email alerts and notifications: - -| Variable | Default | Description | -|----------|---------|-------------| -| `SMTP_HOST` | `""` | SMTP server hostname | -| `SMTP_PORT` | `587` | SMTP server port | -| `SMTP_USER` | `""` | SMTP username | -| `SMTP_PASS` | `""` | SMTP password | -| `SMTP_FROM` | `dequel@localhost` | From address for outgoing emails | - -Config file equivalent: - -```json -{ - "smtpHost": "smtp.sendgrid.net", - "smtpPort": 587, - "smtpUser": "apikey", - "smtpPass": "...", - "smtpFrom": "dequel@example.com" -} -``` - -On boot, these values seed the `smtp_settings` table. The password is encrypted at rest using `ENV_ENCRYPTION_KEY`. You can also update these from the Settings page in the dashboard, and send a test email to verify the configuration. +SMTP settings are not read from environment variables or the config file. Configure them from the Settings page in the dashboard, where you can also send a test email to verify the setup. The password is encrypted at rest using `ENV_ENCRYPTION_KEY`. ### Ingress @@ -98,11 +76,6 @@ Config file equivalent: "githubClientSecret": "...", "githubAppName": "MyDequel", "githubWebhookSecret": "...", - "smtpHost": "smtp.sendgrid.net", - "smtpPort": 587, - "smtpUser": "apikey", - "smtpPass": "...", - "smtpFrom": "alerts@example.com", "envEncryptionKey": "your-secure-key-here" } ``` diff --git a/apps/web/src/components/project/alerts/AlertsTab.tsx b/apps/web/src/components/project/alerts/AlertsTab.tsx index 2c06095..4dec282 100644 --- a/apps/web/src/components/project/alerts/AlertsTab.tsx +++ b/apps/web/src/components/project/alerts/AlertsTab.tsx @@ -1,5 +1,5 @@ import { useQuery } from "@tanstack/react-query"; -import { AlertTriangle, BellOff, Cpu, Layers, Plus, ShieldAlert, Trash2, WifiOff } from "lucide-react"; +import { BellOff, Cpu, Layers, Plus, ShieldAlert, Trash2, WifiOff } from "lucide-react"; import type React from "react"; import { useState } from "react"; import * as api from "../../../api/client"; @@ -18,12 +18,8 @@ const getAlertIcon = (type: string) => { return Cpu; case "memory": return Layers; - case "error_rate": - return AlertTriangle; case "downtime": return WifiOff; - case "cert_expiry": - return ShieldAlert; default: return Cpu; } @@ -33,17 +29,20 @@ const getUnit = (type: string) => { switch (type) { case "cpu": case "memory": - case "error_rate": return "%"; - case "cert_expiry": - return " days"; - case "downtime": - return "s"; default: return ""; } }; +const safeHostname = (url: string) => { + try { + return new URL(url).hostname; + } catch { + return url; + } +}; + export function AlertsTab({ projectId }: AlertsTabProps) { const { data: alerts = [], refetch } = useQuery({ queryKey: ["alerts", projectId], @@ -54,6 +53,7 @@ export function AlertsTab({ projectId }: AlertsTabProps) { const [type, setType] = useState("cpu"); const [threshold, setThreshold] = useState("80"); const [channel, setChannel] = useState("email"); + const [destination, setDestination] = useState(""); const [deletingAlertId, setDeletingAlertId] = useState(null); @@ -68,9 +68,11 @@ export function AlertsTab({ projectId }: AlertsTabProps) { e.preventDefault(); await api.createAlert(projectId, { type, - threshold: Number(threshold), + threshold: type === "downtime" ? null : Number(threshold), channel, + ...(channel !== "email" ? { destination } : {}), } as any); + setDestination(""); setIsOpen(false); refetch(); }; @@ -86,8 +88,7 @@ export function AlertsTab({ projectId }: AlertsTabProps) {

    No Alert Rules Configured

    - Monitor system health and receive notifications when CPU, memory, error rates, or certificate status cross - your limits. + Monitor system health and receive notifications when CPU, memory, or downtime conditions trigger.