From be5afed04dbb83190c46de0e322cb09992dfb611 Mon Sep 17 00:00:00 2001 From: ryzrr Date: Mon, 20 Jul 2026 20:17:53 +0530 Subject: [PATCH] feat: real alert delivery channels (Slack/Discord/PagerDuty/webhook) Backs the Alerts "Channels" tab with real per-project destinations instead of mock state: a new alert_channels table, CRUD + test-delivery endpoints, and an evaluation worker that polls alert rules every 30s and delivers firing/resolved notifications via core/notify.deliver. Also fixes anomaly detection to only alert on metric increases (a drop in error rate or p99 latency is an improvement, not an incident), and shows the active project's real name in the Alerts/Endpoints/Traces breadcrumbs instead of a hardcoded placeholder. --- app/(dashboard)/alerts/page.tsx | 297 +++++++++++------- app/(dashboard)/endpoints/page.tsx | 4 +- app/(dashboard)/traces/page.tsx | 4 +- apps/api/api/routes/alerts.py | 112 +++++++ apps/api/api/schemas.py | 31 ++ apps/api/core/config.py | 2 + apps/api/core/notify.py | 62 ++++ .../migrations/versions/006_alert_channels.py | 50 +++ apps/api/worker/alerts.py | 150 +++++++++ apps/api/worker/anomaly.py | 24 +- apps/api/worker/main.py | 5 +- components/alerts/rule-list.tsx | 5 +- hooks/use-api-query.ts | 11 +- hooks/use-data.ts | 81 ++++- 14 files changed, 709 insertions(+), 129 deletions(-) create mode 100644 apps/api/core/notify.py create mode 100644 apps/api/migrations/versions/006_alert_channels.py create mode 100644 apps/api/worker/alerts.py diff --git a/app/(dashboard)/alerts/page.tsx b/app/(dashboard)/alerts/page.tsx index a7f8a2d..acd47e4 100644 --- a/app/(dashboard)/alerts/page.tsx +++ b/app/(dashboard)/alerts/page.tsx @@ -7,32 +7,27 @@ import { RuleList } from "@/components/alerts/rule-list"; import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; import { cn } from "@/lib/utils"; -import { useAlertRules, useAlertHistory } from "@/hooks/use-data"; -import type { CreateAlertRulePayload } from "@/hooks/use-data"; +import { useAlertRules, useAlertHistory, useChannels } from "@/hooks/use-data"; +import type { CreateAlertRulePayload, AlertChannel, CreateChannelPayload } from "@/hooks/use-data"; +import { useProjects } from "@/components/providers/project-provider"; import { timeAgo } from "@/lib/utils"; // ─── Types ──────────────────────────────────────────────────────────────────── type ChannelType = "slack" | "discord" | "pagerduty" | "webhook"; -interface Channel { - id: string; - type: ChannelType; - name: string; - status: "connected" | "disconnected" | "testing"; - webhookUrl?: string; - lastDelivery?: Date; - lastDeliveryOk?: boolean; - icon: string; -} - -// ─── Mock channel state ─────────────────────────────────────────────────────── +const CHANNEL_ICON: Record = { + slack: "💬", + discord: "🎮", + pagerduty: "📟", + webhook: "🔗", +}; -const INITIAL_CHANNELS: Channel[] = [ - { id: "ch1", type: "slack", name: "Slack #alerts", status: "connected", webhookUrl: "https://hooks.slack.com/...", lastDelivery: new Date(Date.now() - 3600000), lastDeliveryOk: true, icon: "💬" }, - { id: "ch2", type: "pagerduty", name: "PagerDuty", status: "connected", webhookUrl: "https://events.pagerduty.com/...", lastDelivery: new Date(Date.now() - 7200000), lastDeliveryOk: true, icon: "📟" }, - { id: "ch3", type: "discord", name: "Discord #ops", status: "disconnected", icon: "🎮" }, - { id: "ch4", type: "webhook", name: "Custom Webhook", status: "disconnected", icon: "🔗" }, +const CHANNEL_TYPES: { value: ChannelType; label: string }[] = [ + { value: "slack", label: "Slack" }, + { value: "discord", label: "Discord" }, + { value: "pagerduty", label: "PagerDuty" }, + { value: "webhook", label: "Webhook" }, ]; @@ -46,53 +41,61 @@ type Severity = "critical" | "warning" | "info"; // ─── Channel config form ────────────────────────────────────────────────────── -function ChannelConfigForm({ channel, onSave, onClose }: { - channel: Channel; - onSave: (ch: Channel) => void; +const CHANNEL_PLACEHOLDER: Record = { + slack: "https://hooks.slack.com/services/T.../B.../...", + discord: "https://discord.com/api/webhooks/...", + pagerduty: "PagerDuty Events v2 routing key", + webhook: "https://your-server.com/webhook", +}; + +function ChannelConfigForm({ channel, onSave, onTest, onDelete, onClose }: { + channel: AlertChannel; + onSave: (id: string, patch: { webhook_url: string; enabled: boolean }) => Promise; + onTest: (id: string) => Promise<{ ok: boolean; detail: string }>; + onDelete: (id: string) => Promise; onClose: () => void; }) { const [url, setUrl] = useState(channel.webhookUrl ?? ""); const [testing, setTesting] = useState(false); - const [testResult, setTestResult] = useState<"ok" | "fail" | null>(null); + const [saving, setSaving] = useState(false); + const [testResult, setTestResult] = useState<{ ok: boolean; detail: string } | null>(null); - const placeholder: Record = { - slack: "https://hooks.slack.com/services/T.../B.../...", - discord: "https://discord.com/api/webhooks/...", - pagerduty: "https://events.pagerduty.com/v2/enqueue", - webhook: "https://your-server.com/webhook", - }; - - function handleTest() { + async function handleTest() { + setSaving(true); + // Persist the URL first so the backend test uses the latest value. + try { await onSave(channel.id, { webhook_url: url, enabled: true }); } finally { setSaving(false); } setTesting(true); setTestResult(null); - setTimeout(() => { - setTesting(false); - setTestResult("ok"); - }, 1400); + try { setTestResult(await onTest(channel.id)); } + catch { setTestResult({ ok: false, detail: "request failed" }); } + finally { setTesting(false); } } - function handleSave() { - onSave({ ...channel, webhookUrl: url, status: url ? "connected" : "disconnected" }); - onClose(); + async function handleSave() { + setSaving(true); + try { await onSave(channel.id, { webhook_url: url, enabled: true }); onClose(); } + finally { setSaving(false); } } return (
- {channel.icon} + {CHANNEL_ICON[channel.type]} {channel.name}
- + { setUrl(e.target.value); setTestResult(null); }} - placeholder={placeholder[channel.type]} + placeholder={CHANNEL_PLACEHOLDER[channel.type]} className="mt-1 w-full bg-[#111] border border-[#2A2A2A] rounded px-3 py-1.5 text-xs text-[#F5F5F5] placeholder-[#333] outline-none focus:border-blue transition-colors font-mono" />
@@ -100,16 +103,16 @@ function ChannelConfigForm({ channel, onSave, onClose }: { {testResult && (
- {testResult === "ok" ? : } - {testResult === "ok" ? "Test message delivered successfully" : "Delivery failed — check URL"} + {testResult.ok ? : } + {testResult.ok ? `Test message delivered (${testResult.detail})` : `Delivery failed — ${testResult.detail}`}
)} - {channel.lastDelivery && ( + {channel.lastDeliveryAt && (

- Last delivery: {timeAgo(channel.lastDelivery)} ·{" "} + Last delivery: {timeAgo(channel.lastDeliveryAt)} ·{" "} {channel.lastDeliveryOk ? "OK" : "Failed"} @@ -117,18 +120,74 @@ function ChannelConfigForm({ channel, onSave, onClose }: { )}

- - +
); } +function AddChannelForm({ onCreate, onClose }: { + onCreate: (input: CreateChannelPayload) => Promise; + onClose: () => void; +}) { + const [type, setType] = useState("slack"); + const [name, setName] = useState(""); + const [url, setUrl] = useState(""); + const [saving, setSaving] = useState(false); + + async function handleCreate() { + if (!name.trim() || !url.trim()) return; + setSaving(true); + try { await onCreate({ type, name: name.trim(), webhook_url: url.trim(), enabled: true }); onClose(); } + finally { setSaving(false); } + } + + return ( +
+
+ New channel + +
+
+ + setName(e.target.value)} + placeholder="Name (e.g. Slack #alerts)" + className="bg-[#111] border border-[#2A2A2A] rounded px-3 py-1.5 text-xs text-[#F5F5F5] placeholder-[#333] outline-none focus:border-blue flex-1 min-w-[160px]" + /> +
+ setUrl(e.target.value)} + placeholder={CHANNEL_PLACEHOLDER[type]} + className="w-full bg-[#111] border border-[#2A2A2A] rounded px-3 py-1.5 text-xs text-[#F5F5F5] placeholder-[#333] outline-none focus:border-blue font-mono" + /> + +
+ ); +} + // ─── Page ───────────────────────────────────────────────────────────────────── const METRIC_TO_KEY: Record = { @@ -145,8 +204,10 @@ const OPERATOR_TO_SYM: Record" | "<" | "=" | "!="> = { }; export default function AlertsPage() { + const { activeProject } = useProjects(); const { data: fetchedRules, createRule, toggleRule } = useAlertRules(); const { data: historyEntries } = useAlertHistory(); + const { data: channels, createChannel, updateChannel, deleteChannel, testChannel } = useChannels(); // Local rules state — seeded from API, updated optimistically on create const [localRules, setLocalRules] = useState(fetchedRules); @@ -154,8 +215,8 @@ export default function AlertsPage() { useEffect(() => { setLocalRules(fetchedRules); }, [fetchedRules]); const [activeTab, setActiveTab] = useState<"rules" | "channels" | "history">("rules"); - const [channels, setChannels] = useState(INITIAL_CHANNELS); const [configuringId, setConfiguringId] = useState(null); + const [addingChannel, setAddingChannel] = useState(false); const [creating, setCreating] = useState(false); // Rule builder state @@ -164,7 +225,15 @@ export default function AlertsPage() { const [threshold, setThreshold] = useState("5"); const [window, setWindow] = useState("5"); const [severity, setSeverity] = useState("warning"); - const [channel, setChannel] = useState("ch1"); + // Rule "send to" holds the channel NAME (matched by the evaluation worker). + const [channelName, setChannelName] = useState(""); + + // Default the rule-builder channel to the first enabled channel once loaded. + useEffect(() => { + if (!channelName && channels.length > 0) { + setChannelName(channels.find((c) => c.enabled)?.name ?? channels[0].name); + } + }, [channels, channelName]); const firing = localRules.filter((r) => r.status === "firing").length; const triggeredToday = historyEntries.filter((h) => { @@ -174,9 +243,8 @@ export default function AlertsPage() { d.getMonth() === now.getMonth() && d.getDate() === now.getDate(); }).length; - const selectedChannel = channels.find((c) => c.id === channel); - const preview = `Alert when ${metric} ${operator} ${threshold}${metric === "Error Rate" ? "%" : "ms"} for ${window} min → ${selectedChannel?.name ?? "—"} [${severity}]`; + const preview = `Alert when ${metric} ${operator} ${threshold}${metric === "Error Rate" ? "%" : "ms"} for ${window} min → ${channelName || "no channel"} [${severity}]`; async function handleCreateRule() { const payload: CreateAlertRulePayload = { @@ -186,23 +254,12 @@ export default function AlertsPage() { threshold: parseFloat(threshold) || 0, window: parseInt(window, 10) || 5, severity, - channel: selectedChannel?.name ?? "", + channel: channelName, }; setCreating(true); try { const newRule = await createRule(payload); setLocalRules((prev) => [...prev, newRule]); - } catch { - // API unavailable — add optimistic entry so UI reflects intent - setLocalRules((prev) => [ - ...prev, - { - id: `rule_local_${Date.now()}`, - ...payload, - status: "ok" as const, - enabled: true, - }, - ]); } finally { setCreating(false); } @@ -216,14 +273,10 @@ export default function AlertsPage() { } } - function handleSaveChannel(updated: Channel) { - setChannels((prev) => prev.map((c) => c.id === updated.id ? updated : c)); - } - return (
@@ -334,16 +387,19 @@ export default function AlertsPage() { Send to - @@ -361,47 +417,64 @@ export default function AlertsPage() { {activeTab === "channels" && (
- {channels.map((ch) => ( -
-
-
- {ch.icon} -
-

{ch.name}

-
- - {ch.status === "testing" ? "testing…" : ch.status} - - {ch.lastDelivery && ( - - {timeAgo(ch.lastDelivery)} - - )} + {channels.map((ch) => { + const connected = Boolean(ch.webhookUrl) && ch.enabled; + return ( +
+
+
+ {CHANNEL_ICON[ch.type]} +
+

{ch.name}

+
+ + {connected ? "connected" : "not configured"} + + {ch.lastDeliveryAt && ( + + {timeAgo(ch.lastDeliveryAt)} · {ch.lastDeliveryOk ? "OK" : "failed"} + + )} +
+
- + {configuringId === ch.id && ( + { await deleteChannel(id); setConfiguringId(null); }} + onClose={() => setConfiguringId(null)} + /> + )}
- {configuringId === ch.id && ( - setConfiguringId(null)} - /> - )} -
- ))} + ); + })} + {channels.length === 0 && !addingChannel && ( +

+ No channels yet. Add one to receive alert notifications. +

+ )}
- + {addingChannel ? ( + { await createChannel(i); }} onClose={() => setAddingChannel(false)} /> + ) : ( + + )}
)} diff --git a/app/(dashboard)/endpoints/page.tsx b/app/(dashboard)/endpoints/page.tsx index f1f8e6e..2874a2e 100644 --- a/app/(dashboard)/endpoints/page.tsx +++ b/app/(dashboard)/endpoints/page.tsx @@ -11,6 +11,7 @@ import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; import { useTimeRange } from "@/hooks/use-time-range"; import { useEndpoints } from "@/hooks/use-data"; +import { useProjects } from "@/components/providers/project-provider"; import type { Endpoint } from "@/lib/types"; const METHODS = ["ALL", "GET", "POST", "PUT", "DELETE", "PATCH"] as const; @@ -22,6 +23,7 @@ const STATUS_RANGES = [ ] as const; export default function EndpointsPage() { + const { activeProject } = useProjects(); const { range, setRange } = useTimeRange("24h"); const [selected, setSelected] = useState(null); const [search, setSearch] = useState(""); @@ -97,7 +99,7 @@ export default function EndpointsPage() { return (
} /> diff --git a/app/(dashboard)/traces/page.tsx b/app/(dashboard)/traces/page.tsx index cbf9507..f41734c 100644 --- a/app/(dashboard)/traces/page.tsx +++ b/app/(dashboard)/traces/page.tsx @@ -9,6 +9,7 @@ import { ServiceMap } from "@/components/traces/service-map"; import { Badge } from "@/components/ui/badge"; import { cn, formatMs, timeAgo } from "@/lib/utils"; import { useTraces } from "@/hooks/use-data"; +import { useProjects } from "@/components/providers/project-provider"; import type { Trace, TraceSpan } from "@/lib/types"; type DurationFilter = "any" | "<100" | "100-500" | "500-1000" | ">1000"; @@ -34,6 +35,7 @@ function matchesDuration(ms: number, filter: DurationFilter): boolean { type ViewMode = "flamegraph" | "servicemap"; export default function TracesPage() { + const { activeProject } = useProjects(); const { data: traces } = useTraces(); const [selectedTrace, setSelectedTrace] = useState(null); @@ -62,7 +64,7 @@ export default function TracesPage() { return (
- +
{/* Trace list */} diff --git a/apps/api/api/routes/alerts.py b/apps/api/api/routes/alerts.py index 07b5bb8..77cd875 100644 --- a/apps/api/api/routes/alerts.py +++ b/apps/api/api/routes/alerts.py @@ -14,9 +14,14 @@ from api.schemas import ( AlertHistoryOut, AlertRuleOut, + ChannelIn, + ChannelOut, + ChannelTestResponse, + ChannelUpdate, CreateAlertRuleRequest, UpdateAlertRuleRequest, ) +from core.notify import deliver router = APIRouter(prefix="/v1", tags=["alerts"]) logger = logging.getLogger(__name__) @@ -192,3 +197,110 @@ async def get_alert_history( ) for row in rows ] + + +# ─── Alert channels (real delivery targets) ────────────────────────────────── + +def _channel_row_to_out(row: asyncpg.Record) -> ChannelOut: + return ChannelOut( + id=str(row["id"]), + type=row["type"], + name=row["name"], + webhook_url=row["webhook_url"], + enabled=row["enabled"], + last_delivery_at=row["last_delivery_at"].isoformat() if row["last_delivery_at"] else None, + last_delivery_ok=row["last_delivery_ok"], + created_at=row["created_at"].isoformat() if row["created_at"] else "", + ) + + +@router.get("/channels", response_model=list[ChannelOut]) +async def list_channels( + project_id: str = Depends(resolve_project_id), + conn: asyncpg.Connection = Depends(scoped_conn), +) -> list[ChannelOut]: + rows = await conn.fetch( + "SELECT * FROM alert_channels WHERE project_id = $1 ORDER BY created_at", + project_id, + ) + return [_channel_row_to_out(r) for r in rows] + + +@router.post("/channels", response_model=ChannelOut, status_code=201) +async def create_channel( + body: ChannelIn, + project_id: str = Depends(resolve_project_id), + conn: asyncpg.Connection = Depends(get_db), +) -> ChannelOut: + row = await conn.fetchrow( + """ + INSERT INTO alert_channels (project_id, type, name, webhook_url, enabled) + VALUES ($1::uuid, $2, $3, $4, $5) + RETURNING * + """, + project_id, body.type, body.name, body.webhook_url, body.enabled, + ) + return _channel_row_to_out(row) + + +@router.patch("/channels/{channel_id}", response_model=ChannelOut) +async def update_channel( + channel_id: str, + body: ChannelUpdate, + project_id: str = Depends(resolve_project_id), + conn: asyncpg.Connection = Depends(get_db), +) -> ChannelOut: + row = await conn.fetchrow( + """ + UPDATE alert_channels + SET name = COALESCE($3, name), + webhook_url = COALESCE($4, webhook_url), + enabled = COALESCE($5, enabled) + WHERE id = $1::uuid AND project_id = $2::uuid + RETURNING * + """, + channel_id, project_id, body.name, body.webhook_url, body.enabled, + ) + if row is None: + raise HTTPException(status_code=404, detail="Channel not found") + return _channel_row_to_out(row) + + +@router.delete("/channels/{channel_id}", status_code=204, response_model=None, response_class=Response) +async def delete_channel( + channel_id: str, + project_id: str = Depends(resolve_project_id), + conn: asyncpg.Connection = Depends(get_db), +) -> Response: + result = await conn.execute( + "DELETE FROM alert_channels WHERE id = $1::uuid AND project_id = $2::uuid", + channel_id, project_id, + ) + if result == "DELETE 0": + raise HTTPException(status_code=404, detail="Channel not found") + return Response(status_code=204) + + +@router.post("/channels/{channel_id}/test", response_model=ChannelTestResponse) +async def test_channel( + channel_id: str, + project_id: str = Depends(resolve_project_id), + conn: asyncpg.Connection = Depends(get_db), +) -> ChannelTestResponse: + row = await conn.fetchrow( + "SELECT type, webhook_url FROM alert_channels WHERE id = $1::uuid AND project_id = $2::uuid", + channel_id, project_id, + ) + if row is None: + raise HTTPException(status_code=404, detail="Channel not found") + ok, detail = await deliver( + row["type"], row["webhook_url"], + title="Liveboard test alert", + summary="This is a test message from Liveboard. If you can see this, delivery works.", + severity="info", + ) + await conn.execute( + "UPDATE alert_channels SET last_delivery_at = now(), last_delivery_ok = $2 WHERE id = $1::uuid", + channel_id, ok, + ) + return ChannelTestResponse(ok=ok, detail=detail) diff --git a/apps/api/api/schemas.py b/apps/api/api/schemas.py index 1dc52d2..4e99d36 100644 --- a/apps/api/api/schemas.py +++ b/apps/api/api/schemas.py @@ -188,3 +188,34 @@ class AlertHistoryOut(BaseModel): channel: str resolved: bool duration: str + + +# ─── Alert channels (Phase: real alerting) ──────────────────────────────────── + +class ChannelIn(BaseModel): + type: str = Field(..., pattern="^(slack|discord|pagerduty|webhook)$") + name: str = Field(..., min_length=1, max_length=100) + webhook_url: Optional[str] = Field(default=None, max_length=1000) + enabled: bool = True + + +class ChannelUpdate(BaseModel): + name: Optional[str] = Field(default=None, max_length=100) + webhook_url: Optional[str] = Field(default=None, max_length=1000) + enabled: Optional[bool] = None + + +class ChannelOut(BaseModel): + id: str + type: str + name: str + webhook_url: Optional[str] + enabled: bool + last_delivery_at: Optional[str] + last_delivery_ok: Optional[bool] + created_at: str + + +class ChannelTestResponse(BaseModel): + ok: bool + detail: str diff --git a/apps/api/core/config.py b/apps/api/core/config.py index f44e66d..047aae4 100644 --- a/apps/api/core/config.py +++ b/apps/api/core/config.py @@ -25,6 +25,8 @@ class Settings(BaseSettings): # ── AI (Phase 5 — Cerebras) ─────────────────────────────────────────────── cerebras_api_key: SecretStr = SecretStr("") + # Model id must be one your key can access (see `client.models.list()`). + cerebras_model: str = "gpt-oss-120b" # ── CORS / realtime origins (Phase 8.5) ────────────────────────────────── # Comma-separated list of allowed browser origins (dashboard + status page). diff --git a/apps/api/core/notify.py b/apps/api/core/notify.py new file mode 100644 index 0000000..e60fbdc --- /dev/null +++ b/apps/api/core/notify.py @@ -0,0 +1,62 @@ +""" +Alert notification delivery — real HTTP POST to the configured channel. + +Slack / Discord / generic webhooks are just JSON POSTs to an incoming-webhook +URL. PagerDuty uses its Events API v2 (webhook_url holds the routing key). +Returns (ok, detail) and never raises — a delivery failure must not crash the +evaluation loop. +""" +from __future__ import annotations + +import logging + +import httpx + +logger = logging.getLogger(__name__) + +_TIMEOUT = 8.0 + + +async def deliver( + channel_type: str, + webhook_url: str | None, + *, + title: str, + summary: str, + severity: str, +) -> tuple[bool, str]: + if not webhook_url: + return False, "no webhook url configured" + + try: + async with httpx.AsyncClient(timeout=_TIMEOUT) as client: + if channel_type == "slack": + r = await client.post(webhook_url, json={"text": f":rotating_light: *{title}*\n{summary}"}) + elif channel_type == "discord": + r = await client.post(webhook_url, json={"content": f"🚨 **{title}**\n{summary}"}) + elif channel_type == "pagerduty": + r = await client.post( + "https://events.pagerduty.com/v2/enqueue", + json={ + "routing_key": webhook_url, + "event_action": "trigger", + "payload": { + "summary": f"{title} — {summary}"[:1024], + "severity": "critical" if severity == "critical" else "warning", + "source": "liveboard", + }, + }, + ) + else: # generic webhook + r = await client.post( + webhook_url, + json={"title": title, "summary": summary, "severity": severity, "source": "liveboard"}, + ) + + ok = 200 <= r.status_code < 300 + if not ok: + logger.warning("Delivery to %s failed: HTTP %s", channel_type, r.status_code) + return ok, f"HTTP {r.status_code}" + except Exception as exc: # network / timeout / bad url + logger.warning("Delivery to %s errored: %s", channel_type, exc) + return False, str(exc)[:160] diff --git a/apps/api/migrations/versions/006_alert_channels.py b/apps/api/migrations/versions/006_alert_channels.py new file mode 100644 index 0000000..9534c41 --- /dev/null +++ b/apps/api/migrations/versions/006_alert_channels.py @@ -0,0 +1,50 @@ +"""Alert delivery channels (real Slack/Discord/PagerDuty/webhook targets) + +Revision ID: 006 +Revises: 005 +Create Date: 2026-07-19 + +Backs the alerts "Channels" tab with real, per-project destinations. Alert +rules reference a channel by name (alert_rules.channel); the evaluation worker +looks the channel up here to get its webhook URL and delivers the notification. +""" + +from alembic import op + +revision = "006" +down_revision = "005" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.execute(""" + CREATE TABLE IF NOT EXISTS alert_channels ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + project_id UUID NOT NULL REFERENCES projects(id) ON DELETE CASCADE, + type TEXT NOT NULL CHECK (type IN ('slack', 'discord', 'pagerduty', 'webhook')), + name TEXT NOT NULL, + webhook_url TEXT, + enabled BOOLEAN NOT NULL DEFAULT TRUE, + last_delivery_at TIMESTAMPTZ, + last_delivery_ok BOOLEAN, + created_at TIMESTAMPTZ DEFAULT now() + ); + """) + op.execute("CREATE INDEX IF NOT EXISTS idx_alert_channels_project ON alert_channels(project_id);") + + # RLS backstop — dashboard reads (dashboard_reader role) are scoped to the + # active project via app.project_id, same as the other tenant tables. + op.execute("GRANT SELECT ON alert_channels TO dashboard_reader;") + op.execute("ALTER TABLE alert_channels ENABLE ROW LEVEL SECURITY;") + op.execute(""" + CREATE POLICY alert_channels_tenant_isolation ON alert_channels + FOR SELECT + TO dashboard_reader + USING (project_id = NULLIF(current_setting('app.project_id', true), '')::uuid); + """) + + +def downgrade() -> None: + op.execute("DROP POLICY IF EXISTS alert_channels_tenant_isolation ON alert_channels;") + op.execute("DROP TABLE IF EXISTS alert_channels;") diff --git a/apps/api/worker/alerts.py b/apps/api/worker/alerts.py new file mode 100644 index 0000000..a685361 --- /dev/null +++ b/apps/api/worker/alerts.py @@ -0,0 +1,150 @@ +""" +Alert-rule evaluation worker. + +Every 30 s, for each ENABLED rule: + 1. Compute its metric over the rule's window from the events table. + 2. Compare against the threshold using the rule's operator. + 3. State machine: + - breach + not already firing -> status='firing', last_triggered=now, + insert alert_history row, deliver to the rule's channel + - no breach + currently firing -> status='ok', resolve the open history row + Owner DB connection (bypasses RLS); the worker sees every project's rules. +""" +from __future__ import annotations + +import asyncio +import logging + +from core.database import get_pool +from core.notify import deliver + +logger = logging.getLogger(__name__) + +_INTERVAL = 30 # seconds + +_METRIC = { + "error_rate": ("Error rate", "%"), + "p99_latency": ("p99 latency", "ms"), + "requests_per_min": ("Requests/min", ""), +} +_OP_WORD = {">": "exceeds", "<": "drops below", "=": "equals", "!=": "differs from"} + + +def _compare(value: float, operator: str, threshold: float) -> bool: + if operator == ">": + return value > threshold + if operator == "<": + return value < threshold + if operator == "=": + return abs(value - threshold) < 1e-6 + if operator == "!=": + return abs(value - threshold) >= 1e-6 + return False + + +async def _metric_value(conn, project_id, metric: str, window_min: int) -> float | None: + """Return the current metric value over the window, or None if not evaluable.""" + row = await conn.fetchrow( + """ + SELECT + COUNT(*)::int AS total, + COALESCE(100.0 * SUM((status_code >= 400)::int) / NULLIF(COUNT(*), 0), 0)::float AS error_rate, + COALESCE(percentile_cont(0.99) WITHIN GROUP (ORDER BY duration_ms), 0)::float AS p99 + FROM events + WHERE project_id = $1 AND time > now() - ($2 * interval '1 minute') + """, + project_id, + window_min, + ) + total = row["total"] or 0 + if metric == "requests_per_min": + return total / max(window_min, 1) + if total == 0: + return None # can't judge error rate / latency with no traffic + if metric == "error_rate": + return float(row["error_rate"]) + if metric == "p99_latency": + return float(row["p99"]) + return None # unsupported metric (e.g. apdex) — skip + + +async def _deliver_for_rule(conn, project_id, rule, title: str, summary: str) -> None: + ch = await conn.fetchrow( + """ + SELECT id, type, webhook_url FROM alert_channels + WHERE project_id = $1 AND name = $2 AND enabled = TRUE + LIMIT 1 + """, + project_id, + rule["channel"], + ) + if ch is None: + return + ok, _detail = await deliver(ch["type"], ch["webhook_url"], title=title, summary=summary, severity=rule["severity"]) + await conn.execute( + "UPDATE alert_channels SET last_delivery_at = now(), last_delivery_ok = $2 WHERE id = $1", + ch["id"], + ok, + ) + + +async def _evaluate_all() -> None: + pool = await get_pool() + async with pool.acquire() as conn: + rules = await conn.fetch( + """ + SELECT id, project_id, name, metric, operator, threshold, window_min, + severity, channel, status + FROM alert_rules + WHERE enabled = TRUE + """ + ) + for rule in rules: + value = await _metric_value(conn, rule["project_id"], rule["metric"], rule["window_min"]) + if value is None: + continue + + breach = _compare(value, rule["operator"], rule["threshold"]) + label, unit = _METRIC.get(rule["metric"], (rule["metric"], "")) + op_word = _OP_WORD.get(rule["operator"], rule["operator"]) + + if breach and rule["status"] != "firing": + title = f"{label} {op_word} {rule['threshold']:g}{unit}" + summary = ( + f"{label} hit {value:.1f}{unit} over the last {rule['window_min']}m " + f"(threshold {rule['operator']} {rule['threshold']:g}{unit})." + ) + await conn.execute( + "UPDATE alert_rules SET status='firing', last_triggered=now(), updated_at=now() WHERE id=$1", + rule["id"], + ) + await conn.execute( + """ + INSERT INTO alert_history (project_id, rule_id, rule_name, fired_at, channel) + VALUES ($1, $2, $3, now(), $4) + """, + rule["project_id"], rule["id"], rule["name"], rule["channel"], + ) + await _deliver_for_rule(conn, rule["project_id"], rule, title, summary) + logger.info("Alert FIRING rule=%s project=%s value=%.1f", rule["name"], rule["project_id"], value) + + elif not breach and rule["status"] == "firing": + await conn.execute( + "UPDATE alert_rules SET status='ok', updated_at=now() WHERE id=$1", + rule["id"], + ) + await conn.execute( + "UPDATE alert_history SET resolved_at=now() WHERE rule_id=$1 AND resolved_at IS NULL", + rule["id"], + ) + logger.info("Alert RESOLVED rule=%s project=%s", rule["name"], rule["project_id"]) + + +async def run() -> None: + logger.info("Alert evaluator started — checking every %d s", _INTERVAL) + while True: + try: + await _evaluate_all() + except Exception: + logger.exception("Alert evaluation cycle failed") + await asyncio.sleep(_INTERVAL) diff --git a/apps/api/worker/anomaly.py b/apps/api/worker/anomaly.py index 49c0f8a..1117ab6 100644 --- a/apps/api/worker/anomaly.py +++ b/apps/api/worker/anomaly.py @@ -125,7 +125,10 @@ async def _check_project(project_id: str, conn, redis) -> None: z_score = (current - mean) / stddev - if abs(z_score) <= _ZSCORE_THRESHOLD: + # Only INCREASES matter for these "higher = worse" metrics. A drop in + # error rate or p99 latency is an improvement, not an incident — never + # page someone because things got better. + if z_score <= _ZSCORE_THRESHOLD: continue logger.info( @@ -198,7 +201,7 @@ async def _handle_anomaly( ) # 5. Build incident fields - severity = "critical" if abs(z_score) > 4.0 else "warning" + severity = "critical" if z_score > 4.0 else "warning" title = _anomaly_title(metric, current, z_score, unit) top_route = ( f"{context_rows[0]['method']} {context_rows[0]['route']}" @@ -263,24 +266,27 @@ async def _generate_summary( ] or [" No recent traffic data available"] prompt = ( - f"An API anomaly was detected. Write a 2-3 sentence plain-English incident summary.\n\n" + f"You are an SRE assistant. An API anomaly was detected. Write a tight " + f"incident briefing for the on-call engineer.\n\n" f"Anomaly:\n" f" Metric: {metric_label}\n" f" Current value: {current:.1f}{unit}\n" f" Normal range: {mean:.1f} ± {stddev:.1f}{unit}\n" - f" Z-score: {z_score:.1f} (threshold ±3.0)\n\n" + f" Z-score: {z_score:.1f} (threshold +3.0; the metric spiked ABOVE normal)\n\n" f"Recent traffic (last 10 min):\n" + "\n".join(context_lines) - + "\n\nBe concise and direct. Identify what changed, which endpoints are " - "affected, and the likely cause if visible from the data. No markdown, " - "no bullet points, no headers." + + "\n\nIn 3-4 sentences, plain prose (no markdown, no bullets, no headers): " + "(1) state what changed and by how much, (2) name the specific endpoint(s) " + "most likely responsible from the traffic above, (3) give the most likely " + "cause, and (4) recommend one concrete next step the engineer should take " + "right now." ) try: response = await client.chat.completions.create( - model="llama-3.3-70b", + model=settings.cerebras_model, messages=[{"role": "user", "content": prompt}], - max_tokens=200, + max_tokens=260, ) return response.choices[0].message.content.strip() except Exception as exc: diff --git a/apps/api/worker/main.py b/apps/api/worker/main.py index 70ce22a..19ed504 100644 --- a/apps/api/worker/main.py +++ b/apps/api/worker/main.py @@ -14,6 +14,7 @@ import signal from worker.aggregator import run as run_aggregator +from worker.alerts import run as run_alerts from worker.anomaly import run as run_anomaly from worker.metrics import run as run_metrics @@ -29,18 +30,20 @@ async def main() -> None: agg_task = asyncio.create_task(run_aggregator(), name="aggregator") metrics_task = asyncio.create_task(run_metrics(), name="metrics-publisher") anomaly_task = asyncio.create_task(run_anomaly(), name="anomaly-detector") + alerts_task = asyncio.create_task(run_alerts(), name="alert-evaluator") def _shutdown(sig: signal.Signals) -> None: logger.info("Received %s, stopping worker", sig.name) agg_task.cancel() metrics_task.cancel() anomaly_task.cancel() + alerts_task.cancel() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, _shutdown, sig) try: - await asyncio.gather(agg_task, metrics_task, anomaly_task, return_exceptions=True) + await asyncio.gather(agg_task, metrics_task, anomaly_task, alerts_task, return_exceptions=True) except asyncio.CancelledError: pass diff --git a/components/alerts/rule-list.tsx b/components/alerts/rule-list.tsx index 9b88b2f..4457de1 100644 --- a/components/alerts/rule-list.tsx +++ b/components/alerts/rule-list.tsx @@ -1,6 +1,6 @@ "use client"; -import { useState } from "react"; +import { useEffect, useState } from "react"; import { Zap, Bell, Info } from "lucide-react"; import { Badge } from "@/components/ui/badge"; import { timeAgo, cn } from "@/lib/utils"; @@ -33,7 +33,10 @@ const SEVERITY_ICONS = { }; export function RuleList({ rules, onToggle }: RuleListProps) { + // Optimistic copy for instant toggle feedback — re-synced whenever the real + // rules prop changes (e.g. project switch, or empty when there are none). const [localRules, setLocalRules] = useState(rules); + useEffect(() => { setLocalRules(rules); }, [rules]); const toggle = (id: string) => { const rule = localRules.find((r) => r.id === id); diff --git a/hooks/use-api-query.ts b/hooks/use-api-query.ts index 103b5b9..ca4d2c6 100644 --- a/hooks/use-api-query.ts +++ b/hooks/use-api-query.ts @@ -4,6 +4,11 @@ import { useEffect, useRef, useState } from "react"; import { apiFetch } from "@/lib/api-client"; import { useProjects } from "@/components/providers/project-provider"; +// Stable empty-array reference for the "real project, no data" state — reused +// across renders so consumers depending on the result don't see a new [] each +// render (which would loop any effect keyed on it). +const EMPTY_LIST: readonly never[] = []; + /** * Generic hook for REST API data. * - Shows `fallback` (mock) while loading or on error. @@ -90,8 +95,10 @@ export function useApiQuery( // Honesty rule: a REAL project never shows the demo `fallback`. While loading // or on error it shows an empty result (empty state), not fabricated data. // The demo `fallback` is only used when there's no project (landing / loading). - const emptyFallback = (Array.isArray(fallback) ? ([] as unknown as T) : fallback); - const effectiveFallback = projectId ? emptyFallback : fallback; + // Use the shared EMPTY_LIST so the reference is stable across renders. + const effectiveFallback = projectId + ? (Array.isArray(fallback) ? (EMPTY_LIST as unknown as T) : fallback) + : fallback; return { data: data ?? effectiveFallback, loading, refetch }; } diff --git a/hooks/use-data.ts b/hooks/use-data.ts index 420fd2c..316123e 100644 --- a/hooks/use-data.ts +++ b/hooks/use-data.ts @@ -283,7 +283,8 @@ export function useAlertRules() { const { activeProject } = useProjects(); const project = activeProject?.id; const fallback = useMemo(() => getAlertRules(), []); - const query = useApiQuery("/v1/alert-rules", {}, fallback, toAlertRules); + // Poll so a rule flipping to "firing" (by the evaluation worker) shows up live. + const query = useApiQuery("/v1/alert-rules", {}, fallback, toAlertRules, 15_000); async function createRule(payload: CreateAlertRulePayload): Promise { const raw = await apiMutate("POST", "/v1/alert-rules", payload, project); @@ -304,5 +305,81 @@ export function useAlertRules() { /** Alert firing history for the current project. */ export function useAlertHistory() { const fallback = useMemo((): AlertHistoryEntry[] => [], []); - return useApiQuery("/v1/alert-history", {}, fallback, toAlertHistory); + return useApiQuery("/v1/alert-history", {}, fallback, toAlertHistory, 15_000); +} + +// ─── Alert channels (real delivery targets) ────────────────────────────────── + +export interface AlertChannel { + id: string; + type: "slack" | "discord" | "pagerduty" | "webhook"; + name: string; + webhookUrl: string | null; + enabled: boolean; + lastDeliveryAt: Date | null; + lastDeliveryOk: boolean | null; +} + +interface ApiChannel { + id: string; + type: string; + name: string; + webhook_url: string | null; + enabled: boolean; + last_delivery_at: string | null; + last_delivery_ok: boolean | null; +} + +function toChannel(r: ApiChannel): AlertChannel { + return { + id: r.id, + type: r.type as AlertChannel["type"], + name: r.name, + webhookUrl: r.webhook_url, + enabled: r.enabled, + lastDeliveryAt: r.last_delivery_at ? new Date(r.last_delivery_at) : null, + lastDeliveryOk: r.last_delivery_ok, + }; +} + +export interface CreateChannelPayload { + type: AlertChannel["type"]; + name: string; + webhook_url: string; + enabled?: boolean; +} + +/** Real alert delivery channels for the current project. */ +export function useChannels() { + const { activeProject } = useProjects(); + const project = activeProject?.id; + const fallback = useMemo(() => [], []); + const query = useApiQuery( + "/v1/channels", + {}, + fallback, + (raw) => (raw as ApiChannel[]).map(toChannel), + 15_000, + ); + + async function createChannel(input: CreateChannelPayload): Promise { + const raw = await apiMutate("POST", "/v1/channels", input, project); + query.refetch(); + return toChannel(raw); + } + async function updateChannel(id: string, patch: Partial): Promise { + await apiMutate("PATCH", `/v1/channels/${id}`, patch, project); + query.refetch(); + } + async function deleteChannel(id: string): Promise { + await apiMutate("DELETE", `/v1/channels/${id}`, undefined, project); + query.refetch(); + } + async function testChannel(id: string): Promise<{ ok: boolean; detail: string }> { + const res = await apiMutate<{ ok: boolean; detail: string }>("POST", `/v1/channels/${id}/test`, {}, project); + query.refetch(); + return res; + } + + return { ...query, createChannel, updateChannel, deleteChannel, testChannel }; }