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 }; }