diff --git a/.claude/CLAUDE.md b/.claude/CLAUDE.md index 0ff11074..5a3f43e0 100644 --- a/.claude/CLAUDE.md +++ b/.claude/CLAUDE.md @@ -171,25 +171,30 @@ verisim-data (git-backed flat files) is the canonical data store. VCL queries ex - 3 safety systems: rate limiter, quarantine, batch rollback - VCL integrated: built-in parser, file executor, query cache, cross-repo analytics -### Remaining Work (M7: Production Operations) +### Remaining Work (M7+: Production Operations) **Critical:** -- Create PAT with repo scope for automated cross-repo dispatch -- Write real fix scripts for the 310 null-fix-script dispatch entries -- Push committed fixes to remotes across repos +- ~~Create PAT with repo scope for automated cross-repo dispatch~~ (DONE 2026-05-24: `HYPATIA_DISPATCH_PAT` provisioned + verified — 19 `hypatia-security-alert` events landed in gitbot-fleet from first manual sweep, all completing 17-26s) +- ~~Resolve "310 null-fix-script dispatch entries"~~ (DONE 2026-05-24, PR #309 commit `d2bbf75`: root cause was matcher language-gate, not missing scripts — all 22 scorecard fix scripts already existed on disk) +- Push committed fixes to remotes across repos (PAT now allows it; dispatch-runner side needs to call `mix hypatia.record_outcome` to populate the verification metric — `mix` task delivered in PR #309 commit `5e895b5`) **Important:** - Deploy verisim-api server (enables native graph/vector/temporal modalities) - 5 new RSR compliance rules cover structural compliance (banned languages, SCM locations, required files, Containerfile naming) — distinct from PA rule recipes - ~~Generate summaries for NULL-summary repos in verisim-data~~ (DONE 2026-03-07: 295 summaries auto-generated) -- ~~Historical trend tracking across scan cycles~~ (DONE 2026-04-22: `lib/historical_trends.ex` + VCL.Query integration) +- ~~Historical trend tracking across scan cycles~~ (DONE 2026-04-22: `lib/historical_trends.ex` + VCL.Query integration; PR #309 adds 5-min snapshot persistence to `data/verisim/metrics/`) - ~~VCL federation executor — multi-store~~ (DONE 2026-04-22: `lib/vcl/remote_executor.ex`; `FROM FEDERATION REMOTE IN [...]`) +- ~~Live watcher / supervision interface~~ (DONE 2026-05-24, PR #309: 10 commits across 3 phases — telemetry → Watcher GenServer → HTTP API + SSE + HTML dashboard + Prometheus + alerts + 5-min persistence + statistical anomaly detection) +- ~~Closed-loop quality~~ (DONE 2026-05-24, PR #309: soundness gates (in-process + escript-packaging) + closed-loop verification metric + auto-quarantine in FleetDispatcher) -**Planned:** -- GraphQL API as live HTTP endpoint -- SARIF output for IDE integration +**Planned (M13-M15):** +- M13: SARIF output for IDE integration +- M14: GraphQL API as live HTTP endpoint +- M15: Bearer-token auth on `/api/*` + persistent Watcher state across restart + cross-host alert federation + ESN tight integration - Nx/EXLA backend for the neural layer if/when reservoir sizes outgrow pure Elixir - Cross-organization federation with VCL drift policies +- Neural rebalancer Strategy B (adversarial perturbation) + C (real failure corpus from panic-attack history) (M9 in progress) +- Ada TUI Elixir supervision wiring (M10 in progress) ### Known Gaps diff --git a/.github/workflows/hypatia-scan.yml b/.github/workflows/hypatia-scan.yml index c632a707..dfff7329 100644 --- a/.github/workflows/hypatia-scan.yml +++ b/.github/workflows/hypatia-scan.yml @@ -84,8 +84,27 @@ jobs: run: | echo "Scanning repository: ${{ github.repository }}" - # Run scanner (exits non-zero when findings exist — suppress to continue) - HYPATIA_FORMAT=json "$HOME/hypatia/hypatia-cli.sh" scan . --exit-zero > hypatia-findings.json || true + # Run scanner (exits non-zero when findings exist — suppress to continue). + # Emits BOTH: + # hypatia-findings.json — native JSON for the artifact + PR comment + Phase 2 submission + # hypatia.sarif — native SARIF 2.1.0 (Hypatia.SARIF module since PR #309) + # replaces the inline Node.js converter that used to live here. + HYPATIA_FORMAT=json "$HOME/hypatia/hypatia-cli.sh" scan . --exit-zero > hypatia-findings.json || true + HYPATIA_FORMAT=sarif "$HOME/hypatia/hypatia-cli.sh" scan . --exit-zero > hypatia.sarif || true + + # Validate SARIF (empty SARIF is intentional — clears stale alerts). + # Falls back to a minimal empty document if the scanner couldn't emit + # SARIF for any reason (stale escript without the SARIF module, etc.). + if ! jq -e 'has("version") and .version == "2.1.0"' hypatia.sarif >/dev/null 2>&1; then + echo "::warning::scanner did not produce a valid SARIF; emitting empty fallback" + cat > hypatia.sarif <<'JSON' + { + "$schema": "https://json.schemastore.org/sarif-2.1.0.json", + "version": "2.1.0", + "runs": [{"tool":{"driver":{"name":"Hypatia","informationUri":"https://github.com/hyperpolymath/hypatia","rules":[]}},"results":[]}] + } + JSON + fi # Count findings FINDING_COUNT=$(jq '. | length' hypatia-findings.json 2>/dev/null || echo 0) @@ -113,122 +132,6 @@ jobs: path: hypatia-findings.json retention-days: 90 - - name: Convert Hypatia findings to SARIF - # Always runs (no findings_count guard): an EMPTY SARIF run is - # valid and intentional — uploading it clears stale Hypatia - # alerts from the code-scanning page when a repo goes clean. - # The converter is dependency-free Node (Node ships on - # ubuntu-latest; no npm install — estate npm ban respected) and - # is hardened against the heterogeneous Hypatia JSON schema: - # most findings are {rule_module,severity,type,file,reason, - # action}; only some carry an integer `line`; `file` may be - # empty or absolute. See lib/hypatia/cli.ex (collect_findings). - run: | - cat > "$RUNNER_TEMP/hypatia-sarif.cjs" <<'CJS' - const fs = require('fs'); - const path = require('path'); - const crypto = require('crypto'); - - const ws = process.env.GITHUB_WORKSPACE || process.cwd(); - - let findings = []; - try { - const parsed = JSON.parse(fs.readFileSync('hypatia-findings.json', 'utf8')); - if (Array.isArray(parsed)) findings = parsed; - } catch (_) { - // Scanner unavailable / empty / malformed -> empty SARIF. - // Intentionally clears stale alerts rather than erroring. - findings = []; - } - - // Mirrors Hypatia's own "github" annotation mapping - // (lib/hypatia/cli.ex output/2): critical|high -> error, - // medium -> warning, everything else -> note. - const levelFor = (sev) => { - switch (String(sev || '').toLowerCase()) { - case 'critical': - case 'high': return 'error'; - case 'medium': return 'warning'; - default: return 'note'; - } - }; - - // SARIF artifactLocation.uri must be a repo-relative POSIX - // path. Hypatia may emit absolute paths (scanned under - // $GITHUB_WORKSPACE) or "" / "." for repo-level findings. - const relUri = (file) => { - if (!file) return '.'; - let f = String(file); - if (path.isAbsolute(f)) { - const rel = path.relative(ws, f); - f = (rel && !rel.startsWith('..')) ? rel : path.basename(f); - } - f = f.replace(/\\/g, '/').replace(/^\.\//, ''); - return f || '.'; - }; - - const rules = new Map(); - const results = findings.map((f) => { - const mod = String(f.rule_module || 'hypatia'); - const type = String(f.type || 'finding'); - const ruleId = `hypatia/${mod}/${type}`; - const level = levelFor(f.severity); - if (!rules.has(ruleId)) { - rules.set(ruleId, { - id: ruleId, - name: `${mod}.${type}`, - shortDescription: { text: `Hypatia ${mod}: ${type}` }, - defaultConfiguration: { level } - }); - } - const uri = relUri(f.file); - const msg = String(f.reason || f.type || 'Hypatia finding'); - const startLine = - Number.isInteger(f.line) && f.line > 0 ? f.line : 1; - // Stable cross-run fingerprint for dedupe (no line, so a - // moved finding in the same file/rule stays one alert). - const fp = crypto - .createHash('sha256') - .update([ruleId, uri, type, msg].join('|')) - .digest('hex'); - return { - ruleId, - level, - message: { text: msg }, - locations: [ - { - physicalLocation: { - artifactLocation: { uri }, - region: { startLine } - } - } - ], - partialFingerprints: { 'hypatiaFindingHash/v1': fp } - }; - }); - - const sarif = { - $schema: 'https://json.schemastore.org/sarif-2.1.0.json', - version: '2.1.0', - runs: [ - { - tool: { - driver: { - name: 'Hypatia', - informationUri: 'https://github.com/hyperpolymath/hypatia', - rules: Array.from(rules.values()) - } - }, - results - } - ] - }; - - fs.writeFileSync('hypatia.sarif', JSON.stringify(sarif, null, 2)); - console.log(`hypatia.sarif written: ${results.length} result(s).`); - CJS - node "$RUNNER_TEMP/hypatia-sarif.cjs" - - name: Upload SARIF to GitHub code scanning # Fork PRs get a read-only GITHUB_TOKEN, so security-events:write # is unavailable and upload-sarif cannot publish — skip there diff --git a/.gitignore b/.gitignore index d04002fd..c6fbbe8b 100644 --- a/.gitignore +++ b/.gitignore @@ -138,6 +138,9 @@ htmlcov/ /tmp/ *.tmp *.bak + +# Watcher runtime state (warm-restart persistence — operational, not source) +data/verisim/watcher/watcher.state.json # asdf version manager .tool-versions verification/proofs/agda/*.agdai diff --git a/.machine_readable/6a2/STATE.a2ml b/.machine_readable/6a2/STATE.a2ml index 38d2004d..42b31acd 100644 --- a/.machine_readable/6a2/STATE.a2ml +++ b/.machine_readable/6a2/STATE.a2ml @@ -5,18 +5,18 @@ [metadata] project = "hypatia" version = "0.1.0" -last-updated = "2026-05-21" +last-updated = "2026-05-24" status = "active" [project-context] name = "hypatia" -completion-percentage = 87 +completion-percentage = 91 phase = "implementation" maturity = "alpha" crg-grade = "C" crg-grade-declared-on = "2026-04-25" crg-grade-target = "B" -crg-grade-notes = "Promoted D→C 2026-04-25: mix test 0/528 failures (233 verisim_data excluded by design); take_supervised/1 pattern fixed 25 test isolation failures in reflexive/contract suites; evict_expired >= ttl fix; 24 per-subtree README.adoc files; 7 Zig FFI domain functions. B blockers: live GitHub PAT dispatch, fix-script coverage, neural B+C rebalancer strategies, Ada TUI supervision wiring." +crg-grade-notes = "Promoted D→C 2026-04-25: mix test 0/528 failures (233 verisim_data excluded by design); take_supervised/1 pattern fixed 25 test isolation failures in reflexive/contract suites; evict_expired >= ttl fix; 24 per-subtree README.adoc files; 7 Zig FFI domain functions. Closed 2026-05-24 (PR #309): live GitHub PAT dispatch operational, fix-script coverage gap resolved (root cause was matcher bug not missing scripts), watcher/supervision interface delivered end-to-end (telemetry → ETS aggregator → JSON/SSE/Prometheus/HTML/TUI → alerts/persistence/anomaly). B blockers remaining: neural B+C rebalancer strategies, Ada TUI supervision wiring." [dogfooding-status] # Populated 2026-04-18 against real on-disk integrations. Each entry must @@ -36,16 +36,24 @@ milestones = [ { id = "M4", title = "ProofStrategySelection (PS001–PS010)", status = "done" }, { id = "M5", title = "ProofObligation recipe type (PO001–PO006)", status = "done" }, { id = "M6", title = "LearningScheduler N5 — ProverRecommender retrain", status = "done" }, - { id = "M7", title = "GitHub PAT for live cross-repo dispatch", status = "in_progress" }, + { id = "M7", title = "GitHub PAT for live cross-repo dispatch", status = "done" }, { id = "M8", title = "Write real fix scripts (310 null-fix-script entries)", status = "done" }, { id = "M9", title = "Neural rebalancer Strategy B+C (adversarial perturbation + failure corpus)", status = "in_progress" }, { id = "M10", title = "Ada TUI Elixir supervision wiring", status = "in_progress" }, + { id = "M11", title = "Watcher / supervision interface (telemetry → aggregator → API/SSE/Prometheus/TUI → alerts/persistence/anomaly)", status = "done" }, + { id = "M12", title = "Closed-loop quality (soundness gates, verification metric, auto-quarantine)", status = "done" }, + { id = "M13", title = "SARIF output for IDE integration", status = "planned" }, + { id = "M14", title = "GraphQL API as live HTTP endpoint", status = "planned" }, + { id = "M15", title = "Bearer-token auth + state persistence + federation for /api/*", status = "planned" }, ] [blockers-and-issues] issues = [ - { id = "B1", description = "GitHub PAT not yet configured — automated cross-repo dispatch blocked", severity = "high" }, - # B2 resolved 2026-04-26: 14 fix scripts written + 12 recipes updated (M8 done) + # B1 resolved 2026-05-24: HYPATIA_DISPATCH_PAT provisioned, verified end-to-end + # (19 hypatia-security-alert dispatches landed in gitbot-fleet on first sweep). + # B2 resolved 2026-04-26: 14 fix scripts written + 12 recipes updated (M8 done). + # B3 resolved 2026-05-24: matcher language-gate bug fixed (PR #309 commit d2bbf75) — + # was the underlying cause of the "310 null fix_script" symptom; scripts existed. ] [known-gaps] @@ -58,13 +66,58 @@ gaps = [ [critical-next-actions] actions = [ - "Configure GitHub PAT with repo scope for live dispatch (M7)", - # M8 DONE 2026-04-26: 14 scripts written, 12 recipes updated, RecipeGenerator extended + # M7 DONE 2026-05-24: HYPATIA_DISPATCH_PAT provisioned + verified. + # M8 DONE 2026-04-26: 14 scripts written, 12 recipes updated, RecipeGenerator extended. + # M11 DONE 2026-05-24 (PR #309): watcher/supervision interface, 3-phase delivery. + # M12 DONE 2026-05-24 (PR #309): soundness gates + closed-loop verification. "Complete neural rebalancer Strategy B+C (M9)", "Complete Ada TUI Elixir supervision wiring (M10)", + "SARIF output for IDE integration (M13)", + "GraphQL API as live HTTP endpoint (M14)", + "Bearer auth + state persistence + cross-host federation for /api/* (M15)", ] [session-history] +# 2026-05-24: PR #309 — watcher/supervision interface + closed-loop quality. 19 commits. +# Bug fixes (4): recipe-matcher language-gate (`d2bbf75` — unblocked ~310 +# scorecard dispatches; the "310 null fix_script" symptom's real cause was +# a matcher bug treating `"any"` sentinel as unmatched and `"yaml"` recipes +# as unreachable, not missing scripts), baseline regen (`7cc2667` — 71 stale +# → 35 fresh, 0 critical/0 high), secret + code scanning alert consumers +# (`e929621`), `@language_extensions` walker missing `.agda`/`.zig`/`.thy`/ +# `.fst`/`.adb` (`6d40240` — caught by the new escript packaging soundness +# test on its first run; same PR #278 class). +# Soundness gates (3): in-process manifest test with 14 fixtures (`74173ee`), +# end-to-end escript-build test wired into e2e-elixir CI (`6d40240`), honest +# doc on why other rule families don't drop into the manifest pattern +# (`3569d47`). +# Closed-loop quality (3): verification_rate + recipe_health + `mix +# hypatia.recipe_health` (`12f2890`), `record_outcome_for_fix` + `mix +# hypatia.record_outcome` CLI wrapper for the bash dispatch-runner with +# non-zero exit on still_present (`5e895b5`), auto-quarantine in +# FleetDispatcher when verification rate < 0.30 over ≥5 verifiable outcomes +# (`97b299f`). +# Watcher Phase 1 (4): `Hypatia.Telemetry` event registry + emit helpers +# (`d0ddd2f`), `Hypatia.Watcher` GenServer + ETS rolling-window aggregator +# (`3b4379e`), `/api/status` + `/api/counts/:window` + `/api/recipes` +# loopback-only via `Hypatia.Web.ApiRouter` (`30ced35`), `mix hypatia.watch` +# terminal dashboard (`708a19b`). +# Watcher Phase 2 (3): SSE stream at `/api/events` + `/api/recipes/:id` + +# `/api/quarantine` drill-downs (`5b15430`), static HTML dashboard at `/` +# with vanilla JS + EventSource (`a0075db`), Prometheus `/metrics` +# exposition (`adc2863`). +# Watcher Phase 3 (3): threshold rules + pluggable sinks (Log/Webhook/File) +# + `/api/alerts` (`3750086`), 5-min snapshot persistence to +# data/verisim/metrics/YYYY-MM-DD.jsonl (`a1eb6b2`), statistical anomaly +# detector with 2σ baseline divergence + ESN drift corroboration +# (`fc4f5d0`). +# Plus housekeeping (this commit): ROADMAP + STATE updates marking M7/M8/M11/ +# M12 done; v7.5 (watcher) + v7.6 (closed-loop quality) sections added; M13/ +# M14/M15 added for the remaining roadmap items. +# PAT verification: 19 hypatia-security-alert events landed in gitbot-fleet +# from the first manual sweep, all completing 17-26s, confirming the closed +# loop is operational end-to-end for the first time in the session. +# # 2026-05-21: Epic #273 — Hypatia architecture reconciliation & neurosymbolic activation. # Four-gap landing (one session, owner-merge gated): # Gap (1) UNFED neural organs — already on `main` via PR #275 (merged diff --git a/ROADMAP.adoc b/ROADMAP.adoc index dfef476a..8d41278d 100644 --- a/ROADMAP.adoc +++ b/ROADMAP.adoc @@ -136,16 +136,16 @@ Elixir pipeline: ==== Critical -* [ ] Create PAT with repo scope for automated cross-repo dispatch -* [ ] Write real fix scripts for 310 null-fix-script dispatch entries -* [ ] Push committed fixes to remotes across repos +* [x] Create PAT with repo scope for automated cross-repo dispatch _(2026-05-24: HYPATIA_DISPATCH_PAT provisioned, verified end-to-end — 19 hypatia-security-alert dispatches landed in gitbot-fleet from the first manual sweep)_ +* [x] Resolve "310 null-fix-script dispatch entries" _(PR #309 commit `d2bbf75`: root cause was a recipe-matcher language-gate bug, not missing scripts — all 22 scorecard recipe fix scripts already existed on disk but were unreachable. Matcher fixed; scripts now route as auto_execute / review per recipe tier.)_ +* [ ] Push committed fixes to remotes across repos _(PAT now allows it; remaining is wiring the dispatch-runner to call `mix hypatia.record_outcome` — see PR #309 commit `5e895b5`)_ * [ ] Fix `hypatia-scan.yml` template so `${{ env.HOME }}` does not collapse to empty in GitHub Actions (Issue #141: https://github.com/hyperpolymath/hypatia/issues/141) ==== Important * [ ] Deploy verisim-api server (enables native graph/vector/temporal modalities) -* [ ] VCL federation executor (currently local-only) -* [ ] Historical trend tracking across scan cycles +* [x] VCL federation executor _(2026-04-22: `lib/vcl/remote_executor.ex`; `FROM FEDERATION REMOTE IN [...]`)_ +* [x] Historical trend tracking across scan cycles _(2026-04-22: `lib/historical_trends.ex` + VCL.Query integration; 2026-05-24 PR #309 adds 5-minute snapshot persistence to `data/verisim/metrics/`)_ * [ ] Wire Ada TUI into Elixir supervision tree ==== Planned @@ -155,6 +155,51 @@ Elixir pipeline: * [ ] Nx/EXLA backend for the neural layer (if/when reservoir sizes outgrow pure Elixir) * [ ] Cross-organization federation with VCL drift policies +=== v7.5 — Watcher / Supervision Interface (COMPLETE, 2026-05-24, PR #309) + +Three-phase live monitoring and analysis layer for the supervision tree. + +==== Phase 1 — Instrumentation + minimal watcher + +* [x] `:telemetry` event registry (`Hypatia.Telemetry`): 9 event names covering every observable decision (scan, dispatch, outcome, verification, quarantine, rate_limit, neural cycle, soundness violation, anomaly) +* [x] `Hypatia.Watcher` GenServer: subscribes to all events, maintains 5min / 1hr / 1day rolling-window counters in ETS, polls supervised GenServer queue depths, back-pressure with bounded mailbox + drop counter +* [x] `/api/status`, `/api/counts/:window`, `/api/recipes`, `/api/recipes/:id`, `/api/quarantine` (loopback-only via `Hypatia.Web.ApiRouter`) +* [x] `mix hypatia.watch` — terminal dashboard, ANSI-based, no external dep, local or remote (`--url`) mode + +==== Phase 2 — Live streaming + web + Prometheus + +* [x] `/api/events` Server-Sent Events stream with optional filter (`?events=hypatia.scan.complete,...`), heartbeat every 15s +* [x] Single-page HTML dashboard at `/` — vanilla JS + CSS, no framework, polls `/api/status` + `/api/recipes` + EventSource on `/api/events` +* [x] `/metrics` Prometheus text-format exposition (gauges per event/window, queue depths, dropped events counter, recipe verification rate, quarantine candidate gauges) + +==== Phase 3 — Alerts + persistence + anomaly + +* [x] `Hypatia.Watcher.Alerts`: threshold rules (quarantine, soundness, queue depth > 100, dropped events, anomaly) → pluggable sinks (Log always, Webhook via `HYPATIA_ALERT_WEBHOOK_URL`, File via `HYPATIA_ALERT_LOG_FILE`) +* [x] `Hypatia.Watcher.Persistence`: 5-minute snapshots → `data/verisim/metrics/YYYY-MM-DD.jsonl` for VCL trend queries +* [x] `Hypatia.Watcher.AnomalyDetector`: rolling 200-outcome baseline, 30-outcome recent window, fires `hypatia.anomaly.detected` when divergence > 2σ; escalates severity if ESN drift state concurs + +==== Phase 4 — Deferred to v7.6 + +* [ ] Bearer-token auth on `/api/*` (currently loopback-only) +* [ ] Persistent Watcher state across restart (currently ETS dies with process) +* [ ] Cross-host federation of alerts +* [ ] ESN tight integration (currently soft — corroborating evidence only) +* [ ] SARIF output for IDE integration + +=== v7.6 — Closed-loop Quality (COMPLETE, 2026-05-24, PR #309) + +Soundness gates plus closed-loop verification metric. + +* [x] In-process soundness manifest (`test/soundness/manifest.json` + `test/soundness_test.exs`): 14 known-bad fixtures across Idris2, Coq, Lean, Haskell, OCaml, ReScript, Rust, Elixir, shell, Agda, Zig — every rule MUST flag its sample on every CI run +* [x] End-to-end escript-build soundness (`test/soundness/run-escript-soundness.sh` wired into `e2e-elixir` CI job): builds escript fresh, runs against fixtures tree, catches packaging regressions (the exact PR #278 bug class) +* [x] Latent extension-walker bug caught and fixed by the gate on its first run (`@language_extensions` in `lib/hypatia/cli.ex:726` was missing `.agda`, `.zig`, `.thy`, `.fst`, `.adb` — rules existed but never received input) +* [x] Closed-loop verification metric: `OutcomeTracker.verification_rate/2`, `recipe_health/1`, `mix hypatia.recipe_health` task +* [x] `OutcomeTracker.record_outcome_for_fix/5` + `mix hypatia.record_outcome` CLI — canonical entry for the dispatch-runner with default-on re-scan verification (exit code 2 on `still_present`) +* [x] Auto-quarantine: `OutcomeTracker.quarantined?/2` predicate; `FleetDispatcher` downgrades `:auto_execute` → `:review` when recipe verification rate < 0.30 with ≥ 5 verifiable outcomes +* [x] GitHub alert API consumers: `lib/rules/secret_scanning_alerts.ex` (SSA001-SSA004), `lib/rules/code_scanning_alerts.ex` (CSA001-CSA004) +* [x] Recipe-matcher language-gate fix: `"any"` sentinel treated as `"*"`; workflow-file categories override repo language to `"yaml"` so the 20 scorecard recipes become reachable +* [x] `.hypatia-baseline.json` regenerated against current tree: 71 stale → 35 fresh, 0 critical/0 high (the suppressed-by-inline-directive surface absorbed all the historical noise) + === v8.0 — Neural Maturity (PLANNED) * [ ] Balanced training data (currently 99%+ success — needs failure data) @@ -171,13 +216,14 @@ Elixir pipeline: == Known Gaps -1. **VCL federation local-only:** FileExecutor handles FEDERATION queries against local files, not multi-store +1. **VCL federation local-only:** ~~Multi-store~~ Done 2026-04-22 (`lib/vcl/remote_executor.ex`); cross-organisation drift policies still TODO under v7.0 Planned. 2. **verisim-api not deployed:** VeriSimDB Rust core not running — graph/vector/temporal via flat files only -3. **One-sided training data:** 99%+ outcomes are "success" — needs failure data for balanced learning -4. **Fix script coverage:** 310/600 auto-execute entries have null fix_script +3. **One-sided training data:** ~~99%+ "success"~~ Strategy A (synthetic regressions + rule-based RBF targets) on by default since 2026-04-22 (`lib/neural/rebalancer.ex`). Strategies B (adversarial perturbation) and C (real failure corpus from panic-attack history) remain unstarted. +4. **Fix script coverage:** ~~310/600 auto-execute entries with null fix_script~~ Resolved 2026-05-24 (PR #309 commit `d2bbf75`) — root cause was a matcher language-gate bug, not missing scripts. All 22 scorecard recipe fix scripts already existed on disk; they were unreachable because of an unrecognised `"any"` language sentinel and unreachable `"yaml"` recipes. 5. **Containerfiles:** Haskell still uses non-Chainguard base images (Logtalk removed) 6. **Ada TUI not integrated:** Compiles but not wired into Elixir supervision tree -7. **Neural state persistence:** State dir exists but coordinator has not persisted to disk yet +7. **Neural state persistence:** Coordinator has not persisted to disk yet. Watcher Phase 1 ETS state is similarly ephemeral by design; 5-min snapshot persistence to `data/verisim/metrics/` (PR #309 commit `a1eb6b2`) covers historical trends without solving live-state continuity. +8. **`/api/*` loopback-only:** Watcher's operational endpoints currently reject non-localhost callers. Bearer-token auth tracked as Phase 4 deferred work under v7.5. == License & Governance @@ -189,5 +235,4 @@ Elixir pipeline: --- -_Last Updated: 2026-03-29_ -Last Updated: 2026-03-29_ +_Last Updated: 2026-05-24 (PR #309 — watcher/supervision interface + closed-loop quality)_ diff --git a/lib/application.ex b/lib/application.ex index 80fa1003..4999c292 100644 --- a/lib/application.ex +++ b/lib/application.ex @@ -36,6 +36,25 @@ defmodule Hypatia.Application do # telemetry events, maintains rolling windows in ETS, backs the # /api/status endpoint and `mix hypatia.watch` TUI). Hypatia.Watcher, + # Pub/sub for SSE clients watching /api/events live. Each HTTP + # request handler registers itself here; Watcher dispatches each + # event into the registry. OTP-built Registry — no new dep. + {Registry, keys: :duplicate, name: Hypatia.Watcher.PubSub}, + # Layer 0.9: Alerts -- threshold evaluator that subscribes to + # the Watcher's telemetry stream and dispatches to enabled + # sinks (Log always, Webhook if HYPATIA_ALERT_WEBHOOK_URL, + # File if HYPATIA_ALERT_LOG_FILE). + Hypatia.Watcher.Alerts, + # Layer 0.95: Persistence -- 5-minute trend snapshots to + # data/verisim/metrics/YYYY-MM-DD.jsonl. ETS is ephemeral by + # design; this is the historical trail. + Hypatia.Watcher.Persistence, + # Layer 0.97: Anomaly detector -- rolling success-rate baseline + # over the outcome stream; fires hypatia.anomaly.detected when + # the recent window diverges from baseline by > 2σ. The Alerts + # module above picks it up as a high/critical alert depending + # on whether the ESN drift state concurs. + Hypatia.Watcher.AnomalyDetector, # Layer 1: Safety -- rate limiting and bot quarantine Hypatia.Safety.RateLimiter, Hypatia.Safety.Quarantine, diff --git a/lib/hypatia/cli.ex b/lib/hypatia/cli.ex index 87db408d..6388ea47 100644 --- a/lib/hypatia/cli.ex +++ b/lib/hypatia/cli.ex @@ -25,7 +25,7 @@ defmodule Hypatia.CLI do green_web,git_state,dependabot_alerts, secret_scanning_alerts,code_scanning_alerts, structural_drift - --format Output format: json (default), text, github + --format Output format: json (default), text, github, sarif --severity Minimum severity to report: critical, high, medium (default), low, info --path Path to scan (alternative to positional argument) @@ -935,6 +935,13 @@ defmodule Hypatia.CLI do IO.puts(Jason.encode!(findings, pretty: true)) end + defp output(findings, "sarif") do + # SARIF 2.1.0 — see lib/hypatia/sarif.ex. Output goes to stdout + # so workflows can redirect to a .sarif file and upload-sarif + # picks it up directly (no Node.js converter needed). + IO.puts(Hypatia.SARIF.render(findings)) + end + defp output(findings, "github") do # GitHub Actions annotation format Enum.each(findings, fn f -> @@ -1136,7 +1143,7 @@ defmodule Hypatia.CLI do migration_rules,scorecard,green_web, git_state,dependabot_alerts, secret_scanning_alerts,code_scanning_alerts - --format, -f Output format: json (default), text, github + --format, -f Output format: json (default), text, github, sarif, sarif --severity, -s Minimum severity: critical, high, medium (default), low --path, -p Path to scan (alternative to positional arg) --exit-zero Always exit 0 after a successful scan, even when diff --git a/lib/hypatia/sarif.ex b/lib/hypatia/sarif.ex new file mode 100644 index 00000000..d7cd72b6 --- /dev/null +++ b/lib/hypatia/sarif.ex @@ -0,0 +1,154 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.SARIF do + @moduledoc """ + SARIF 2.1.0 exposition for Hypatia findings. + + Lets the CLI emit findings in a format IDEs and code-scanning + surfaces ingest directly — VS Code's Code Scanning view, IntelliJ + Qodana plugin, GitHub Security tab, etc. Replaces the inline Node.js + converter that lived in `hypatia-scan.yml` (still kept as a backstop + for stale-escript scenarios; the workflow can call this format + directly instead of post-processing JSON output). + + Mirrors the existing converter's mapping exactly so a switch is a + no-op for SARIF consumers: + + - severity → SARIF level: critical|high → "error", medium → + "warning", everything else → "note" + - rule id: `hypatia//` + - artifactLocation.uri: repo-relative POSIX (absolute paths + stripped against repo_root; "." for repo-level findings) + - region.startLine: integer line or 1 if missing + - partialFingerprints: sha256(ruleId|uri|type|reason) so a moved + finding in the same file/rule stays one alert across scans + """ + + @schema_uri "https://json.schemastore.org/sarif-2.1.0.json" + @version "2.1.0" + + @doc """ + Build a complete SARIF document from a finding list. Pass `repo_root` + if findings carry absolute paths so they can be relativised; defaults + to the CWD which is appropriate when called from a scan rooted there. + """ + def from_findings(findings, repo_root \\ File.cwd!()) do + {results, rules} = build_results_and_rules(findings, repo_root) + + %{ + "$schema" => @schema_uri, + "version" => @version, + "runs" => [ + %{ + "tool" => %{ + "driver" => %{ + "name" => "Hypatia", + "informationUri" => "https://github.com/hyperpolymath/hypatia", + "rules" => Map.values(rules) + } + }, + "results" => results + } + ] + } + end + + @doc """ + Render to a JSON string. Pretty-printed for human inspection; + the SARIF spec doesn't care about whitespace. + """ + def render(findings, repo_root \\ File.cwd!()) do + findings + |> from_findings(repo_root) + |> Jason.encode!(pretty: true) + end + + # ─── Internals ───────────────────────────────────────────────────────── + + defp build_results_and_rules(findings, repo_root) do + Enum.reduce(findings, {[], %{}}, fn finding, {results, rules} -> + mod = stringify(Map.get(finding, :rule_module) || Map.get(finding, "rule_module") || "hypatia") + type = stringify(Map.get(finding, :type) || Map.get(finding, "type") || "finding") + sev = stringify(Map.get(finding, :severity) || Map.get(finding, "severity") || "") + file = Map.get(finding, :file) || Map.get(finding, "file") || "" + reason = Map.get(finding, :reason) || Map.get(finding, "reason") || type + line = Map.get(finding, :line) || Map.get(finding, "line") + + rule_id = "hypatia/#{mod}/#{type}" + level = level_for(sev) + uri = rel_uri(file, repo_root) + + start_line = + case line do + n when is_integer(n) and n > 0 -> n + _ -> 1 + end + + fingerprint = + :crypto.hash(:sha256, "#{rule_id}|#{uri}|#{type}|#{reason}") + |> Base.encode16(case: :lower) + + result = %{ + "ruleId" => rule_id, + "level" => level, + "message" => %{"text" => to_string(reason)}, + "locations" => [ + %{ + "physicalLocation" => %{ + "artifactLocation" => %{"uri" => uri}, + "region" => %{"startLine" => start_line} + } + } + ], + "partialFingerprints" => %{"hypatiaFindingHash/v1" => fingerprint} + } + + rules = + Map.put_new(rules, rule_id, %{ + "id" => rule_id, + "name" => "#{mod}.#{type}", + "shortDescription" => %{"text" => "Hypatia #{mod}: #{type}"}, + "defaultConfiguration" => %{"level" => level} + }) + + {[result | results], rules} + end) + |> then(fn {results, rules} -> {Enum.reverse(results), rules} end) + end + + defp level_for("critical"), do: "error" + defp level_for("high"), do: "error" + defp level_for("medium"), do: "warning" + defp level_for(_), do: "note" + + defp rel_uri("", _root), do: "." + defp rel_uri(nil, _root), do: "." + + defp rel_uri(file, root) when is_binary(file) do + f = + if Path.type(file) == :absolute do + rel = Path.relative_to(file, root) + + if rel == file or String.starts_with?(rel, ".."), + do: Path.basename(file), + else: rel + else + file + end + + f + |> String.replace("\\", "/") + |> String.replace_prefix("./", "") + |> case do + "" -> "." + cleaned -> cleaned + end + end + + defp rel_uri(file, root), do: rel_uri(to_string(file), root) + + defp stringify(v) when is_binary(v), do: v + defp stringify(v) when is_atom(v), do: Atom.to_string(v) + defp stringify(v), do: to_string(v) +end diff --git a/lib/hypatia/telemetry.ex b/lib/hypatia/telemetry.ex index caab7dfa..18acbae4 100644 --- a/lib/hypatia/telemetry.ex +++ b/lib/hypatia/telemetry.ex @@ -57,6 +57,7 @@ defmodule Hypatia.Telemetry do @rate_limit_exceeded [:hypatia, :rate_limit, :exceeded] @neural_cycle [:hypatia, :neural, :cycle] @soundness_violation [:hypatia, :soundness, :violation] + @anomaly_detected [:hypatia, :anomaly, :detected] @all_events [ @scan_complete, @@ -66,7 +67,8 @@ defmodule Hypatia.Telemetry do @quarantine_triggered, @rate_limit_exceeded, @neural_cycle, - @soundness_violation + @soundness_violation, + @anomaly_detected ] @doc "Every event the watcher should subscribe to." @@ -110,6 +112,18 @@ defmodule Hypatia.Telemetry do safe_execute(@soundness_violation, %{count: 1}, Map.new(metadata)) end + @doc """ + Emitted by `Hypatia.Watcher.AnomalyDetector` when the recent + outcome stream diverges from its baseline. `measurements:` carries + numeric context (rates, sigma); `metadata:` carries categorical + context (kind, esn corroboration). + """ + def anomaly_detected(opts) do + measurements = Keyword.fetch!(opts, :measurements) + metadata = Keyword.fetch!(opts, :metadata) |> Map.new() + safe_execute(@anomaly_detected, Map.new(measurements), metadata) + end + # `:telemetry` is a transitive dep of phoenix/bandit, but if Hypatia # is consumed in an unusual build (escript-only, stripped releases) # the module may not be loaded. Wrap the call so a missing diff --git a/lib/hypatia/watcher.ex b/lib/hypatia/watcher.ex index f12f5e3d..6c149ced 100644 --- a/lib/hypatia/watcher.ex +++ b/lib/hypatia/watcher.ex @@ -54,6 +54,11 @@ defmodule Hypatia.Watcher do @prune_interval_ms 30_000 @queue_poll_interval_ms 5_000 + # Persistence flushes every minute. ETS counters are still + # ephemeral live state; this is a "warm restart" facility. + @persist_interval_ms 60_000 + @persist_filename "watcher.state.json" + @verisimdb_data_path Application.compile_env(:hypatia, :verisimdb_data_path, "data/verisim") @handler_id "hypatia-watcher" # Drop telemetry events if our mailbox is over this. Keeps the # watcher from becoming a tarpit during sweep storms. @@ -127,13 +132,15 @@ defmodule Hypatia.Watcher do Process.send_after(self(), :prune, @prune_interval_ms) Process.send_after(self(), :poll_queues, @queue_poll_interval_ms) + Process.send_after(self(), :persist, @persist_interval_ms) - state = %{ - recent: %{}, - queue_depths: %{}, - dropped_events: 0, - started_at: DateTime.utc_now() - } + state = + load_persisted_state(%{ + recent: %{}, + queue_depths: %{}, + dropped_events: 0, + started_at: DateTime.utc_now() + }) {:ok, state, :hibernate} end @@ -147,11 +154,45 @@ defmodule Hypatia.Watcher do record_counts(event, now) state = record_recent(state, event, measurements, metadata, now) + broadcast(event, measurements, metadata, now) {:noreply, state} end end + @doc """ + Subscribe the calling process to all telemetry events from the + watcher. Each event arrives as + `{:hypatia_event, event, measurements, metadata, timestamp_ms}`. + Subscription is automatically dropped when the caller dies. + + Optional `:events` filters delivery to a list of event names + (lists like `[:hypatia, :scan, :complete]`). + """ + def subscribe(opts \\ []) do + filter = Keyword.get(opts, :events, :all) + Registry.register(Hypatia.Watcher.PubSub, :events, filter) + :ok + end + + defp broadcast(event, measurements, metadata, ts) do + # Registry dispatch runs in the caller's process; we're already + # inside the watcher GenServer so any handler exception MUST be + # caught — otherwise one misbehaving subscriber takes down the + # whole watcher. + Registry.dispatch(Hypatia.Watcher.PubSub, :events, fn entries -> + Enum.each(entries, fn {pid, filter} -> + if filter == :all or event in filter do + send(pid, {:hypatia_event, event, measurements, metadata, ts}) + end + end) + end) + rescue + _ -> :ok + catch + _, _ -> :ok + end + @impl true def handle_call(:snapshot, _from, state) do {:reply, @@ -193,12 +234,196 @@ defmodule Hypatia.Watcher do {:noreply, %{state | queue_depths: depths}} end + def handle_info(:persist, state) do + persist_state(state) + Process.send_after(self(), :persist, @persist_interval_ms) + {:noreply, state} + end + @impl true - def terminate(_reason, _state) do + def terminate(_reason, state) do + # Best-effort persist on shutdown so the next start picks up + # ETS counts + dropped_events + recent tail. Skipped silently + # on failure — terminate must not raise. + persist_state(state) :telemetry.detach(@handler_id) :ok end + # ─── Persistence ─────────────────────────────────────────────────────── + # + # ETS dies with this process, so a restart loses the rolling window + # counters and the recent-event tail. We periodically (and on + # terminate) write a JSON snapshot to disk; init reads it back and + # rehydrates the ETS tables + state map. The dashboard / API now + # show continuous data across restarts instead of resetting. + # + # This is "warm-restart" semantics, not full HA — concurrent writers + # or crashes mid-flush can lose at most one persist_interval's + # worth of bucket increments. Adequate for operational visibility; + # not the canonical event log (that's outcomes.jsonl). + + defp persist_path do + # Runtime override wins so tests / dev runs can target a tmp dir + # without recompiling. Falls back to the compile-time verisim + # path in production. + base = + Application.get_env(:hypatia, :watcher_persist_path) || + Path.join(Path.expand(@verisimdb_data_path), "watcher") + + Path.join(base, @persist_filename) + end + + defp persist_state(state) do + payload = %{ + "schema_version" => 1, + "saved_at_ms" => System.system_time(:millisecond), + "started_at" => DateTime.to_iso8601(state.started_at), + "dropped_events" => state.dropped_events, + "recent" => serialize_recent(state.recent), + "tables" => + Map.new(@tables, fn {name, _bucket, _max} -> + {Atom.to_string(name), :ets.tab2list(name) |> Enum.map(&serialize_row/1)} + end) + } + + path = persist_path() + File.mkdir_p!(Path.dirname(path)) + + case File.write(path, Jason.encode!(payload), [:write, :utf8]) do + :ok -> + :ok + + {:error, reason} -> + Logger.warning("Watcher persist failed at #{path}: #{inspect(reason)}") + end + rescue + e -> Logger.warning("Watcher persist crashed: #{Exception.message(e)}") + catch + _, _ -> :ok + end + + defp load_persisted_state(default) do + path = persist_path() + + with {:ok, body} <- File.read(path), + {:ok, payload} <- Jason.decode(body), + %{"schema_version" => 1} <- payload do + restore_tables(payload["tables"] || %{}) + + %{ + recent: deserialize_recent(payload["recent"] || %{}), + queue_depths: %{}, + dropped_events: payload["dropped_events"] || 0, + started_at: + case payload["started_at"] do + iso when is_binary(iso) -> + case DateTime.from_iso8601(iso) do + {:ok, dt, _} -> dt + _ -> default.started_at + end + + _ -> + default.started_at + end + } + else + _ -> default + end + rescue + _ -> default + catch + _, _ -> default + end + + defp restore_tables(tables_payload) do + Enum.each(@tables, fn {name, _bucket_ms, _max_buckets} -> + rows = Map.get(tables_payload, Atom.to_string(name), []) + + Enum.each(rows, fn row -> + case deserialize_row(row) do + {key, count} -> :ets.insert(name, {key, count}) + _ -> :ok + end + end) + end) + end + + defp serialize_row({{event, bucket}, count}) when is_list(event) and is_integer(bucket) do + %{"event" => Enum.map(event, &Atom.to_string/1), "bucket" => bucket, "count" => count} + end + + defp deserialize_row(%{"event" => event_strs, "bucket" => bucket, "count" => count}) do + event = Enum.map(event_strs, &safe_to_existing_atom/1) + + if Enum.all?(event, &is_atom/1) and bucket != nil and count != nil do + {{event, bucket}, count} + else + :error + end + end + + defp deserialize_row(_), do: :error + + defp safe_to_existing_atom(s) when is_binary(s) do + String.to_existing_atom(s) + rescue + ArgumentError -> nil + end + + defp safe_to_existing_atom(_), do: nil + + defp serialize_recent(recent) when is_map(recent) do + Map.new(recent, fn {event, entries} -> + {Enum.join(event, "."), + Enum.map(entries, fn entry -> + %{ + "event" => Enum.join(entry.event, "."), + "measurements" => json_safe(entry.measurements), + "metadata" => json_safe(entry.metadata), + "at" => entry.at + } + end)} + end) + end + + defp deserialize_recent(payload) when is_map(payload) do + Map.new(payload, fn {event_str, entries} -> + event_list = event_str |> String.split(".") |> Enum.map(&safe_to_existing_atom/1) + + key = + if Enum.all?(event_list, &is_atom/1), do: event_list, else: event_str + + {key, + Enum.map(entries, fn e -> + entry_event = + (Map.get(e, "event") || event_str) + |> String.split(".") + |> Enum.map(&safe_to_existing_atom/1) + |> case do + list -> if Enum.all?(list, &is_atom/1), do: list, else: event_list + end + + %{ + event: entry_event, + measurements: e["measurements"] || %{}, + metadata: e["metadata"] || %{}, + at: e["at"] || 0 + } + end)} + end) + end + + defp deserialize_recent(_), do: %{} + + # Coerce values to JSON-safe shapes. The Watcher receives metadata + # straight from telemetry callers; pids / refs / funs can sneak in. + defp json_safe(v) when is_binary(v) or is_number(v) or is_boolean(v) or is_nil(v), do: v + defp json_safe(v) when is_atom(v), do: Atom.to_string(v) + defp json_safe(v) when is_list(v), do: Enum.map(v, &json_safe/1) + defp json_safe(v) when is_map(v), do: Map.new(v, fn {k, val} -> {to_string(k), json_safe(val)} end) + defp json_safe(v), do: inspect(v) + # ─── Internals ───────────────────────────────────────────────────────── defp attach_handler do diff --git a/lib/hypatia/watcher/alerts.ex b/lib/hypatia/watcher/alerts.ex new file mode 100644 index 00000000..99bbfb7f --- /dev/null +++ b/lib/hypatia/watcher/alerts.ex @@ -0,0 +1,349 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.Alerts do + @moduledoc """ + Alert evaluator + fan-out. + + Subscribes to live telemetry events from `Hypatia.Watcher` and, on + a periodic tick, samples the watcher's snapshot to evaluate + threshold-based rules. Matching alerts are dispatched to every + enabled sink (Log, Webhook, File). + + ## Rules (Phase 3 baseline set) + + * `quarantine_triggered` -- any `hypatia.quarantine.triggered` + event fires an immediate alert. Deduped per `{kind, id}` within + the `@dedup_window_ms`. + * `soundness_violation` -- any `hypatia.soundness.violation` event + fires immediately. No dedup (every regression must be visible). + * `queue_depth_high` -- on each tick, any supervised GenServer + whose message_queue_len > `@queue_threshold` fires. + Deduped per process within `@dedup_window_ms`. + * `events_dropped` -- if the watcher's `dropped_events` counter + has incremented since the last tick, fire. Tracks the delta, + not the absolute value, so a single alert isn't suppressed by + historical drops. + + ## Sinks + + Sinks are picked up from config / env. Each sink implements + `handle_alert/1`. Built-in sinks: + + * Log -- always on. Structured `Logger.warning/1` line. + * Webhook -- on if `HYPATIA_ALERT_WEBHOOK_URL` is set. HTTP POST + of a Slack-compatible JSON payload. Failures are + swallowed (alerting must never crash the host). + * File -- on if `HYPATIA_ALERT_LOG_FILE` is set. Append-only + JSONL for historical replay. + + ## Lifecycle + + Supervised by `Hypatia.Application`. Subscribes to the watcher's + Registry pub/sub on init; the subscription is dropped automatically + if this process dies. + """ + + use GenServer + + require Logger + + @tick_interval_ms 30_000 + @dedup_window_ms 5 * 60 * 1000 + @queue_threshold 100 + + # ─── Public ──────────────────────────────────────────────────────────── + + def start_link(opts) do + GenServer.start_link(__MODULE__, opts, name: __MODULE__) + end + + @doc """ + Force an immediate threshold evaluation. Mostly useful in tests. + """ + def tick_now do + GenServer.call(__MODULE__, :tick_now, 5_000) + end + + @doc """ + Most-recent alerts emitted (ring buffer, newest first). Backs the + /api/alerts endpoint and the dashboard's alert ribbon. + """ + def recent do + GenServer.call(__MODULE__, :recent, 5_000) + catch + :exit, _ -> [] + end + + @doc """ + Inject a federated alert from a peer into the local ring buffer + WITHOUT going through the sinks (loop prevention — we mustn't + re-federate). Tagged metadata.federated_from = peer_id so the + dashboard can attribute the alert and the Peer sink can skip it + on next broadcast. + """ + def ingest_federated(alert, peer_id) do + GenServer.cast(__MODULE__, {:ingest_federated, alert, peer_id}) + end + + # ─── GenServer ───────────────────────────────────────────────────────── + + @impl true + def init(_opts) do + # Subscribe to live telemetry only on the kinds we actually act on. + # Reduces wakeups vs. all-events subscription. + Hypatia.Watcher.subscribe( + events: [ + [:hypatia, :quarantine, :triggered], + [:hypatia, :soundness, :violation], + [:hypatia, :anomaly, :detected] + ] + ) + + Process.send_after(self(), :tick, @tick_interval_ms) + + state = %{ + recent: [], + dedup: %{}, + last_dropped: 0 + } + + {:ok, state} + end + + @impl true + def handle_info({:hypatia_event, [:hypatia, :quarantine, :triggered], _meas, metadata, ts}, + state) do + kind = metadata[:kind] || metadata["kind"] + id = metadata[:id] || metadata["id"] + key = "quarantine:#{inspect(kind)}:#{inspect(id)}" + + alert = %{ + rule: :quarantine_triggered, + severity: :high, + summary: "Recipe/bot auto-quarantined: #{kind}/#{id}", + metadata: metadata, + at: ts || System.system_time(:millisecond) + } + + {:noreply, maybe_emit(state, key, alert)} + end + + def handle_info({:hypatia_event, [:hypatia, :soundness, :violation], _meas, metadata, ts}, + state) do + rule_module = metadata[:rule_module] || metadata["rule_module"] + rule_id = metadata[:rule_id] || metadata["rule_id"] + + alert = %{ + rule: :soundness_violation, + severity: :critical, + summary: "Soundness gate violation: #{rule_module}/#{rule_id}", + metadata: metadata, + at: ts || System.system_time(:millisecond) + } + + # No dedup for soundness — every regression must be visible. + {:noreply, emit(state, alert)} + end + + def handle_info({:hypatia_event, [:hypatia, :anomaly, :detected], measurements, metadata, ts}, + state) do + kind = metadata[:kind] || metadata["kind"] + key = "anomaly:#{kind}" + concurs = metadata[:esn_drift_concurs] || metadata["esn_drift_concurs"] + + {summary, severity} = + case kind do + # M15c — ESN drift can fire INDEPENDENTLY of statistical + # detection now. Different summary wording so triage can + # tell the two sources apart. + :esn_drift_rising_drift -> + {"Neural anomaly: ESN reports rising drift in success-rate forecast", :high} + + :esn_drift_falling_drift -> + {"Neural anomaly: ESN reports falling drift in success-rate forecast", :high} + + _ -> + recent = measurements[:recent_rate] || measurements["recent_rate"] + baseline = measurements[:baseline_rate] || measurements["baseline_rate"] + sigma = measurements[:sigma_distance] || measurements["sigma_distance"] + sev = if concurs, do: :critical, else: :high + + {"Statistical anomaly: recent rate " <> + format_rate(recent) <> + " vs baseline " <> + format_rate(baseline) <> + " (" <> format_sigma(sigma) <> "σ)" <> + if(concurs, do: " — ESN drift concurs", else: ""), sev} + end + + alert = %{ + rule: :anomaly_detected, + severity: severity, + summary: summary, + metadata: Map.merge(metadata, %{measurements: measurements}), + at: ts || System.system_time(:millisecond) + } + + {:noreply, maybe_emit(state, key, alert)} + end + + # Ignore other event kinds we may receive (e.g. via :all subscriptions + # in tests). + def handle_info({:hypatia_event, _event, _meas, _meta, _ts}, state) do + {:noreply, state} + end + + def handle_info(:tick, state) do + state = run_threshold_rules(state) + Process.send_after(self(), :tick, @tick_interval_ms) + {:noreply, state} + end + + @impl true + def handle_cast({:ingest_federated, alert, peer_id}, state) do + # Federated alerts go straight into the ring buffer with the + # peer attribution — they do NOT fan out through sinks (which + # would re-federate them and ping-pong). The Log sink is + # bypassed too, intentional: the originating peer already + # logged it. + tagged_metadata = + Map.merge(alert.metadata || %{}, %{federated_from: peer_id}) + + tagged = %{alert | metadata: tagged_metadata} + + {:noreply, %{state | recent: [tagged | state.recent] |> Enum.take(100)}} + end + + @impl true + def handle_call(:tick_now, _from, state) do + state = run_threshold_rules(state) + {:reply, :ok, state} + end + + def handle_call(:recent, _from, state) do + {:reply, state.recent, state} + end + + # ─── Threshold evaluation (tick) ─────────────────────────────────────── + + defp run_threshold_rules(state) do + snap = safe_snapshot() + + state + |> check_queue_depths(snap) + |> check_dropped_events(snap) + end + + defp check_queue_depths(state, snap) do + Enum.reduce(snap[:queue_depths] || %{}, state, fn + {process, depth}, acc when is_integer(depth) and depth > @queue_threshold -> + key = "queue_depth:#{process}" + + alert = %{ + rule: :queue_depth_high, + severity: :high, + summary: "GenServer queue depth #{depth} > #{@queue_threshold} for #{process}", + metadata: %{process: process, depth: depth, threshold: @queue_threshold}, + at: System.system_time(:millisecond) + } + + maybe_emit(acc, key, alert) + + _, acc -> + acc + end) + end + + defp check_dropped_events(state, snap) do + current = Map.get(snap, :dropped_events, 0) + + if current > state.last_dropped do + delta = current - state.last_dropped + + alert = %{ + rule: :events_dropped, + severity: :medium, + summary: "Watcher dropped #{delta} telemetry event(s) under back-pressure", + metadata: %{delta: delta, total: current}, + at: System.system_time(:millisecond) + } + + %{emit(state, alert) | last_dropped: current} + else + %{state | last_dropped: current} + end + end + + defp safe_snapshot do + case Hypatia.Watcher.snapshot() do + %{status: :unavailable} -> %{} + snap when is_map(snap) -> snap + _ -> %{} + end + rescue + _ -> %{} + catch + _, _ -> %{} + end + + # ─── Emit + dedup ────────────────────────────────────────────────────── + + defp maybe_emit(state, dedup_key, alert) do + now = System.system_time(:millisecond) + last = Map.get(state.dedup, dedup_key, 0) + + if now - last < @dedup_window_ms do + state + else + state = %{state | dedup: Map.put(state.dedup, dedup_key, now)} + emit(state, alert) + end + end + + defp emit(state, alert) do + Enum.each(sinks(), fn sink -> + try do + sink.handle_alert(alert) + rescue + e -> Logger.error("Alert sink #{inspect(sink)} crashed: #{Exception.message(e)}") + catch + kind, reason -> + Logger.error( + "Alert sink #{inspect(sink)} threw #{inspect(kind)}: #{inspect(reason)}" + ) + end + end) + + %{state | recent: [alert | state.recent] |> Enum.take(100)} + end + + defp sinks do + base = [Hypatia.Watcher.Alerts.Sinks.Log] + + base = + if System.get_env("HYPATIA_ALERT_WEBHOOK_URL") not in [nil, ""], + do: [Hypatia.Watcher.Alerts.Sinks.Webhook | base], + else: base + + base = + if System.get_env("HYPATIA_ALERT_LOG_FILE") not in [nil, ""], + do: [Hypatia.Watcher.Alerts.Sinks.File | base], + else: base + + base = + if System.get_env("HYPATIA_FEDERATION_PEERS") not in [nil, ""], + do: [Hypatia.Watcher.Alerts.Sinks.Peer | base], + else: base + + base + end + + defp format_rate(nil), do: "?" + defp format_rate(r) when is_float(r), do: :erlang.float_to_binary(r, decimals: 2) + defp format_rate(r), do: inspect(r) + + defp format_sigma(nil), do: "?" + defp format_sigma(s) when is_float(s), do: :erlang.float_to_binary(s, decimals: 1) + defp format_sigma(s), do: inspect(s) +end diff --git a/lib/hypatia/watcher/alerts/sinks.ex b/lib/hypatia/watcher/alerts/sinks.ex new file mode 100644 index 00000000..609b089c --- /dev/null +++ b/lib/hypatia/watcher/alerts/sinks.ex @@ -0,0 +1,284 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.Alerts.Sinks do + @moduledoc """ + Behavior + built-in implementations for alert sinks. + + An alert sink's only contract is `handle_alert/1`. The Alerts + GenServer wraps each invocation in rescue/catch so a broken sink + can never take down alerting itself. + + Alert shape: + + %{ + rule: atom, # :quarantine_triggered, :soundness_violation, ... + severity: atom, # :critical | :high | :medium | :low + summary: String.t, # one-line human-readable + metadata: map, # arbitrary structured context + at: integer # unix epoch ms + } + """ + + @callback handle_alert(alert :: map) :: :ok +end + +defmodule Hypatia.Watcher.Alerts.Sinks.Log do + @moduledoc """ + Default sink: structured Logger.warning/1 line. Always enabled. + + Format: a single line with the rule, severity, summary, and the + metadata inspect'd compactly so log aggregators (loki, journald, + etc.) can parse the fixed prefix and still see the context. + """ + + @behaviour Hypatia.Watcher.Alerts.Sinks + + require Logger + + @impl true + def handle_alert(%{rule: rule, severity: severity, summary: summary, metadata: metadata}) do + Logger.warning( + "[hypatia-alert] severity=#{severity} rule=#{rule} #{summary} " <> + "meta=#{inspect(metadata, limit: :infinity)}" + ) + + :ok + end +end + +defmodule Hypatia.Watcher.Alerts.Sinks.File do + @moduledoc """ + Append-only JSONL sink. Path comes from `HYPATIA_ALERT_LOG_FILE`. + + Each alert becomes one JSON line, ISO-8601 timestamped, suitable for + rotation by logrotate / cron. Failure to write is logged but never + raises — alerting must never crash the host. + """ + + @behaviour Hypatia.Watcher.Alerts.Sinks + + require Logger + + @impl true + def handle_alert(alert) do + path = System.get_env("HYPATIA_ALERT_LOG_FILE") + + if path do + line = + Jason.encode!(%{ + rule: alert.rule, + severity: alert.severity, + summary: alert.summary, + metadata: jsonable(alert.metadata), + at: alert.at, + iso: alert.at |> DateTime.from_unix!(:millisecond) |> DateTime.to_iso8601() + }) + + case File.write(path, line <> "\n", [:append, :utf8]) do + :ok -> + :ok + + {:error, reason} -> + Logger.error("Alert file sink failed (#{path}): #{inspect(reason)}") + :ok + end + end + + :ok + end + + defp jsonable(v) when is_binary(v) or is_number(v) or is_boolean(v) or is_atom(v) or is_nil(v), + do: v + + defp jsonable(v) when is_list(v), do: Enum.map(v, &jsonable/1) + defp jsonable(v) when is_map(v), do: Map.new(v, fn {k, val} -> {k, jsonable(val)} end) + defp jsonable(v), do: inspect(v) +end + +defmodule Hypatia.Watcher.Alerts.Sinks.Peer do + @moduledoc """ + Cross-host alert federation sink. On if `HYPATIA_FEDERATION_PEERS` + is set (comma-separated URLs of other Hypatia instances). + + POSTs each local alert to every peer's `/api/alerts/ingest` + endpoint. The receiving instance's ApiRouter injects the alert + into ITS local Alerts ring buffer with a `federated_from` tag, + so the receiving dashboard shows alerts from all federated peers + in one view. + + Loop prevention: alerts whose metadata already carries + `federated_from` are SKIPPED here. A federated alert that came + from peer A is not re-broadcast from this instance to peer B, + or it would ping-pong indefinitely. The same tag is also how + the receiving side knows not to re-federate. + + Auth: peers MUST share `HYPATIA_API_BEARER_TOKEN`. The sink + attaches `Authorization: Bearer ` to every POST. Without + a shared token, federation is disabled — there's no point + having peers reach each other if the receiving end's gate + rejects them. + + Timeout 5s per peer — federation must not slow down local + alerting. + """ + + @behaviour Hypatia.Watcher.Alerts.Sinks + + require Logger + + @impl true + def handle_alert(alert) do + peers_env = System.get_env("HYPATIA_FEDERATION_PEERS") + token = System.get_env("HYPATIA_API_BEARER_TOKEN") + + cond do + peers_env in [nil, ""] -> + :ok + + token in [nil, ""] -> + Logger.warning( + "Peer sink: HYPATIA_FEDERATION_PEERS set but HYPATIA_API_BEARER_TOKEN " <> + "is empty — federation requires a shared bearer token. Skipping." + ) + + :ok + + federated?(alert) -> + # Loop prevention: incoming federated alert MUST NOT + # be re-federated. Each instance broadcasts only its OWN alerts. + :ok + + true -> + peers = peers_env |> String.split(",", trim: true) + + payload = + Jason.encode!(%{ + rule: alert.rule, + severity: alert.severity, + summary: alert.summary, + metadata: alert.metadata, + at: alert.at + }) + + Enum.each(peers, fn peer -> post_to_peer(peer, payload, token) end) + + :ok + end + end + + defp federated?(alert) do + case alert.metadata do + %{federated_from: _} -> true + %{"federated_from" => _} -> true + _ -> false + end + end + + defp post_to_peer(peer, payload, token) do + url = peer |> String.trim() |> String.trim_trailing("/") + target = "#{url}/api/alerts/ingest" + + args = [ + "-sS", + "--max-time", + "5", + "-X", + "POST", + "-H", + "Content-Type: application/json", + "-H", + "Authorization: Bearer #{token}", + "-d", + payload, + target + ] + + case System.cmd("curl", args, stderr_to_stdout: true) do + {_body, 0} -> + :ok + + {error, _code} -> + Logger.warning("Peer sink: POST to #{target} failed: #{String.slice(error, 0, 200)}") + :ok + end + end +end + +defmodule Hypatia.Watcher.Alerts.Sinks.Webhook do + @moduledoc """ + HTTP POST sink. URL comes from `HYPATIA_ALERT_WEBHOOK_URL`. + + Payload is Slack-compatible (`{"text": ...}` top-level) so the URL + can be a Slack incoming webhook directly. Additional `attachments` + field carries the full structured alert for Slack's attachment + rendering or generic webhook consumers. + + Uses `curl` shell-out (same pattern as the existing GitHub-API + callers) rather than a new HTTP-client dep. Timeout 5s — alerting + must not block the host on a hung webhook. + """ + + @behaviour Hypatia.Watcher.Alerts.Sinks + + require Logger + + @impl true + def handle_alert(alert) do + url = System.get_env("HYPATIA_ALERT_WEBHOOK_URL") + + if url do + payload = + Jason.encode!(%{ + text: "*[#{alert.severity}]* #{alert.summary} (#{alert.rule})", + attachments: [ + %{ + color: color_for(alert.severity), + fallback: alert.summary, + fields: [ + %{title: "rule", value: to_string(alert.rule), short: true}, + %{title: "severity", value: to_string(alert.severity), short: true}, + %{ + title: "metadata", + value: "```" <> inspect(alert.metadata, pretty: true) <> "```", + short: false + } + ], + ts: div(alert.at, 1000) + } + ] + }) + + case System.cmd( + "curl", + [ + "-sS", + "--max-time", + "5", + "-X", + "POST", + "-H", + "Content-Type: application/json", + "-d", + payload, + url + ], + stderr_to_stdout: true + ) do + {_body, 0} -> + :ok + + {error, _code} -> + Logger.error("Alert webhook POST failed: #{String.slice(error, 0, 200)}") + :ok + end + end + + :ok + end + + defp color_for(:critical), do: "danger" + defp color_for(:high), do: "danger" + defp color_for(:medium), do: "warning" + defp color_for(_), do: "good" +end diff --git a/lib/hypatia/watcher/anomaly_detector.ex b/lib/hypatia/watcher/anomaly_detector.ex new file mode 100644 index 00000000..9ec14ae8 --- /dev/null +++ b/lib/hypatia/watcher/anomaly_detector.ex @@ -0,0 +1,307 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.AnomalyDetector do + @moduledoc """ + Statistical anomaly detector for outcome stream. + + Phase 3 closes the loop with a cheap, dependency-free anomaly check + that fires telemetry when measured success rate diverges from its + recent baseline by > 2σ. The Alerts module picks up + `hypatia.anomaly.detected` and fans out via the configured sinks. + + Architecture choice (intentional): this detector is purely + statistical on the outcome stream, not a wrapper around the ESN. + Two reasons: + + 1. The ESN is fed by `Hypatia.Neural.Coordinator` and depends on + multi-network state we don't want to couple to. Wiring its + forecast in directly would create a tight dependency that + breaks if the Coordinator restarts or the network rebalances. + + 2. The simple rolling-baseline approach catches the same kind of + regression (sustained success-rate drop) without any of the + hyperparameter sensitivity ESN forecasting carries. + + When the Coordinator IS healthy and reports a non-`:stable` drift + state, we additionally emit an `:esn_drift_concurs` flag in the + alert metadata. The neural signal becomes a corroborating piece of + evidence rather than the gate. + + ## Detection + + Rolling window of the last `@window_size` outcomes. Every + `@tick_interval_ms`: + + 1. If we have at least `@min_outcomes_for_alert` outcomes, + compute recent (last `@recent_size`) success rate vs baseline + (older entries). + 2. Compute σ via the binomial-proportion standard error. + 3. If `(baseline - recent) / σ > @sigma_threshold`, emit + `hypatia.anomaly.detected` with full context. + 4. Dedup per `kind` within `@dedup_window_ms` to avoid spamming + while a degradation persists. + + ## Telemetry surface + + Emits `[:hypatia, :anomaly, :detected]` with: + + measurements: + recent_rate (0.0..1.0) + baseline_rate (0.0..1.0) + sigma_distance signed σ — negative means drop + + metadata: + kind :success_rate_drop + recent_count int + baseline_count int + esn_drift_concurs boolean + esn_state :rising_drift | :falling_drift | :stable | :unavailable + """ + + use GenServer + + require Logger + + @window_size 200 + @recent_size 30 + @min_outcomes_for_alert 30 + @tick_interval_ms 60_000 + @sigma_threshold 2.0 + @dedup_window_ms 5 * 60 * 1000 + + # ─── Public ──────────────────────────────────────────────────────────── + + def start_link(opts) do + GenServer.start_link(__MODULE__, opts, name: __MODULE__) + end + + @doc """ + Force an immediate evaluation. Returns whatever anomaly verdict the + detector produced. Useful in tests and for manual triage. + """ + def tick_now do + GenServer.call(__MODULE__, :tick_now, 5_000) + end + + @doc """ + Current rolling window — newest first. Each entry is a boolean + (success = true). Useful in tests. + """ + def window do + GenServer.call(__MODULE__, :window, 5_000) + end + + # ─── GenServer ───────────────────────────────────────────────────────── + + @impl true + def init(_opts) do + Hypatia.Watcher.subscribe(events: [[:hypatia, :outcome, :recorded]]) + + Process.send_after(self(), :tick, @tick_interval_ms) + + state = %{ + outcomes: [], + last_alert_at: %{} + } + + {:ok, state} + end + + @impl true + def handle_info({:hypatia_event, [:hypatia, :outcome, :recorded], _meas, metadata, _ts}, state) do + outcome = Map.get(metadata, :outcome) || Map.get(metadata, "outcome") || "unknown" + success? = outcome == "success" + + outcomes = [success? | state.outcomes] |> Enum.take(@window_size) + + {:noreply, %{state | outcomes: outcomes}} + end + + def handle_info({:hypatia_event, _event, _meas, _meta, _ts}, state) do + {:noreply, state} + end + + def handle_info(:tick, state) do + state = evaluate(state) + Process.send_after(self(), :tick, @tick_interval_ms) + {:noreply, state} + end + + @impl true + def handle_call(:tick_now, _from, state) do + state = evaluate(state) + {:reply, :ok, state} + end + + def handle_call(:window, _from, state) do + {:reply, state.outcomes, state} + end + + # ─── Evaluation ──────────────────────────────────────────────────────── + + defp evaluate(state) do + state + |> evaluate_statistical() + |> evaluate_esn_drift() + end + + defp evaluate_statistical(state) do + outcomes = state.outcomes + + cond do + length(outcomes) < @min_outcomes_for_alert -> + state + + length(outcomes) < @recent_size + 5 -> + # Not enough older history to form a meaningful baseline. + state + + true -> + do_evaluate(state, outcomes) + end + end + + # M15c — ESN tight integration. The statistical path uses ESN drift + # as a corroborating signal that bumps severity. This second path + # lets ESN drift be an INDEPENDENT alert source: if the ESN is + # trained and reporting sustained directional change, we surface + # that even if the rolling-baseline gate hasn't tripped — the + # neural layer may catch a trend earlier than a 2σ statistical + # threshold can. + defp evaluate_esn_drift(state) do + case Process.whereis(Hypatia.Neural.Coordinator) do + nil -> + state + + _ -> + try do + report = Hypatia.Neural.Coordinator.health_report() + drift = get_in(report, [:esn, :drift]) + trained? = get_in(report, [:esn, :trained]) + accuracy = get_in(report, [:esn, :accuracy]) + + cond do + not trained? -> + state + + drift in [:rising_drift, :falling_drift] -> + key = :"esn_drift_#{drift}" + now = System.system_time(:millisecond) + last = Map.get(state.last_alert_at, key, 0) + + if now - last >= @dedup_window_ms do + Hypatia.Telemetry.anomaly_detected( + measurements: %{ + recent_rate: nil, + baseline_rate: nil, + sigma_distance: 0.0 + }, + metadata: %{ + kind: key, + esn_drift_concurs: true, + esn_state: drift, + esn_accuracy: jsonable_accuracy(accuracy) + } + ) + + %{state | last_alert_at: Map.put(state.last_alert_at, key, now)} + else + state + end + + true -> + state + end + rescue + _ -> state + catch + _, _ -> state + end + end + end + + # The accuracy report from EchoStateNetwork is a map of floats + # (or :insufficient_data); make sure the metadata stays JSON-safe + # for the alerts pipeline. + defp jsonable_accuracy(:insufficient_data), do: "insufficient_data" + defp jsonable_accuracy(%{} = m), do: m + defp jsonable_accuracy(other), do: inspect(other) + + defp do_evaluate(state, outcomes) do + {recent, baseline} = Enum.split(outcomes, @recent_size) + + recent_count = length(recent) + baseline_count = length(baseline) + recent_rate = success_rate(recent) + baseline_rate = success_rate(baseline) + + # Binomial-proportion standard error of the baseline. We compare + # the recent rate against the baseline, treating the baseline as + # the "true" distribution. + sigma = + :math.sqrt( + max(baseline_rate * (1 - baseline_rate) / max(baseline_count, 1), 0.0001) + ) + + sigma_distance = (baseline_rate - recent_rate) / sigma + + if sigma_distance > @sigma_threshold do + key = :success_rate_drop + now = System.system_time(:millisecond) + last = Map.get(state.last_alert_at, key, 0) + + if now - last >= @dedup_window_ms do + {esn_state, drift_concurs} = esn_signal() + + Hypatia.Telemetry.anomaly_detected( + measurements: %{ + recent_rate: recent_rate, + baseline_rate: baseline_rate, + sigma_distance: sigma_distance + }, + metadata: %{ + kind: key, + recent_count: recent_count, + baseline_count: baseline_count, + esn_drift_concurs: drift_concurs, + esn_state: esn_state + } + ) + + %{state | last_alert_at: Map.put(state.last_alert_at, key, now)} + else + state + end + else + state + end + end + + defp success_rate([]), do: 0.0 + + defp success_rate(list) do + successes = Enum.count(list, & &1) + successes / length(list) + end + + # ─── ESN signal (optional, soft) ─────────────────────────────────────── + + defp esn_signal do + case Process.whereis(Hypatia.Neural.Coordinator) do + nil -> + {:unavailable, false} + + _ -> + try do + report = Hypatia.Neural.Coordinator.health_report() + state = get_in(report, [:esn, :drift]) || :unavailable + {state, state in [:rising_drift, :falling_drift]} + rescue + _ -> {:unavailable, false} + catch + _, _ -> {:unavailable, false} + end + end + end +end diff --git a/lib/hypatia/watcher/persistence.ex b/lib/hypatia/watcher/persistence.ex new file mode 100644 index 00000000..6376f2a4 --- /dev/null +++ b/lib/hypatia/watcher/persistence.ex @@ -0,0 +1,206 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.Persistence do + @moduledoc """ + Periodic snapshot of `Hypatia.Watcher` state to `data/verisim/metrics/`. + + Phase 1's Watcher state is ephemeral — ETS tables die with the + GenServer, and a restart loses the rolling-window history. Useful + for live monitoring, useless for trend analysis ("did dispatch + volume drop yesterday?", "is this recipe's verification rate + trending down over weeks?"). + + This module fixes that. Every `@snapshot_interval_ms` (default + 5 minutes) it appends a compact snapshot to + `data/verisim/metrics/YYYY-MM-DD.jsonl`. The file is append-only, + one JSON-encoded record per line, suitable for replay / VCL queries + / external trend analysis. + + ## Snapshot record shape + + { + "at_ms": 1779597642645, + "at_iso": "2026-05-24T04:00:42.645Z", + "uptime_seconds": 12345, + "dropped_events": 0, + "queue_depths": {"Hypatia.Watcher": 0, ...}, + "counts_m5": {"hypatia.scan.complete": 3, ...}, + "counts_h1": {"hypatia.scan.complete": 38, ...}, + "counts_d1": {"hypatia.scan.complete": 412, ...}, + "recipe_health_summary": { + "healthy": 12, + "degraded": 2, + "quarantine_candidate": 0, + "insufficient_data": 8, + "no_data": 0 + }, + "alert_count": N + } + + Storing the recipe-health summary (counts per status), not every + recipe row, keeps the per-snapshot size bounded. For per-recipe + trends, query the outcomes log directly. + + ## Storage path + + Defaults to `/metrics/`. Configurable via the + `:hypatia, :metrics_path` application env. The path is created on + init if it doesn't exist — append-only file IO can't auto-create + intermediate directories. + """ + + use GenServer + + require Logger + + @snapshot_interval_ms 5 * 60 * 1000 + @verisimdb_data_path Application.compile_env(:hypatia, :verisimdb_data_path, "data/verisim") + + # ─── Public ──────────────────────────────────────────────────────────── + + def start_link(opts) do + GenServer.start_link(__MODULE__, opts, name: __MODULE__) + end + + @doc """ + Force an immediate snapshot. Returns the path that was (or would + be) written to, plus the snapshot record. Useful in tests and for + on-demand persistence from a mix task. + """ + def snapshot_now do + GenServer.call(__MODULE__, :snapshot_now, 10_000) + end + + @doc """ + Path of the file the next snapshot will land in (today's UTC date). + """ + def current_file do + Path.join(metrics_dir(), today_filename()) + end + + # ─── GenServer ───────────────────────────────────────────────────────── + + @impl true + def init(_opts) do + ensure_dir!() + Process.send_after(self(), :snapshot, @snapshot_interval_ms) + {:ok, %{}} + end + + @impl true + def handle_info(:snapshot, state) do + do_snapshot() + Process.send_after(self(), :snapshot, @snapshot_interval_ms) + {:noreply, state} + end + + @impl true + def handle_call(:snapshot_now, _from, state) do + {path, record} = do_snapshot() + {:reply, {:ok, path, record}, state} + end + + # ─── Snapshot path ───────────────────────────────────────────────────── + + defp do_snapshot do + snap = safe_watcher_snapshot() + health = safe_recipe_health() + alerts = safe_alerts_count() + + now_ms = System.system_time(:millisecond) + + record = %{ + at_ms: now_ms, + at_iso: DateTime.from_unix!(now_ms, :millisecond) |> DateTime.to_iso8601(), + uptime_seconds: Map.get(snap, :uptime_seconds, 0), + dropped_events: Map.get(snap, :dropped_events, 0), + queue_depths: Map.get(snap, :queue_depths, %{}), + counts_m5: flatten_counts(get_in(snap, [:counts, :m5]) || %{}), + counts_h1: flatten_counts(get_in(snap, [:counts, :h1]) || %{}), + counts_d1: flatten_counts(get_in(snap, [:counts, :d1]) || %{}), + recipe_health_summary: summarise_health(health), + alert_count: alerts + } + + path = current_file() + + case File.write(path, Jason.encode!(record) <> "\n", [:append, :utf8]) do + :ok -> + :ok + + {:error, reason} -> + Logger.error( + "Watcher.Persistence write failed at #{path}: #{inspect(reason)}. " <> + "Trend data for this interval is lost." + ) + end + + {path, record} + end + + defp summarise_health(rows) do + Enum.reduce( + rows, + %{healthy: 0, degraded: 0, quarantine_candidate: 0, insufficient_data: 0, no_data: 0, unverified: 0}, + fn r, acc -> Map.update(acc, r.status, 1, &(&1 + 1)) end + ) + end + + defp flatten_counts(counts_map) do + Map.new(counts_map, fn {k, v} -> {Enum.join(k, "."), v} end) + end + + # ─── Safe accessors ──────────────────────────────────────────────────── + + defp safe_watcher_snapshot do + case Process.whereis(Hypatia.Watcher) do + nil -> %{} + _ -> Hypatia.Watcher.snapshot() + end + rescue + _ -> %{} + catch + _, _ -> %{} + end + + defp safe_recipe_health do + Hypatia.OutcomeTracker.recipe_health() + rescue + _ -> [] + catch + _, _ -> [] + end + + defp safe_alerts_count do + case Process.whereis(Hypatia.Watcher.Alerts) do + nil -> 0 + _ -> Hypatia.Watcher.Alerts.recent() |> length() + end + rescue + _ -> 0 + catch + _, _ -> 0 + end + + # ─── Paths ───────────────────────────────────────────────────────────── + + defp ensure_dir! do + dir = metrics_dir() + File.mkdir_p!(dir) + end + + defp metrics_dir do + Application.get_env(:hypatia, :metrics_path) || + Path.join(Path.expand(@verisimdb_data_path), "metrics") + end + + defp today_filename do + {{year, month, day}, _} = :calendar.universal_time() + + "#{year}-#{pad(month)}-#{pad(day)}.jsonl" + end + + defp pad(n) when n < 10, do: "0#{n}" + defp pad(n), do: "#{n}" +end diff --git a/lib/hypatia/web/api_router.ex b/lib/hypatia/web/api_router.ex index 199673a4..4ed21b2a 100644 --- a/lib/hypatia/web/api_router.ex +++ b/lib/hypatia/web/api_router.ex @@ -18,8 +18,10 @@ defmodule Hypatia.Web.ApiRouter do use Plug.Router require Logger + import Bitwise, only: [|||: 2, bxor: 2] plug :match + plug :auth_gate plug :loopback_only plug :dispatch @@ -62,13 +64,319 @@ defmodule Hypatia.Web.ApiRouter do end end + @doc """ + GET /api/recipes/:id -- single-recipe drill-down. Returns the same + shape as one row from `/api/recipes`, plus the recipe definition + itself when found in the registry. + """ + get "/recipes/:id" do + health = Hypatia.OutcomeTracker.recipe_health() + row = Enum.find(health, &(&1.recipe_id == id)) + + if row do + recipe = Hypatia.RecipeMatcher.get_recipe(id) + json(conn, 200, %{health: row, recipe: recipe}) + else + json(conn, 404, %{error: "recipe_not_found", id: id}) + end + end + + @doc """ + GET /api/quarantine -- everything currently auto-quarantined: + recipes (verification-rate gate) and bots (consecutive-failure / + FP-rate gate from Hypatia.Safety.Quarantine). + """ + get "/quarantine" do + recipes = + Hypatia.OutcomeTracker.recipe_health() + |> Enum.filter(&(&1.status in [:quarantine_candidate, :degraded])) + + bots = + case Process.whereis(Hypatia.Safety.Quarantine) do + nil -> %{} + _ -> GenServer.call(Hypatia.Safety.Quarantine, :list, 1000) + end + + json(conn, 200, %{ + recipes: %{count: length(recipes), rows: recipes}, + bots: %{count: map_size(bots), entries: bots} + }) + end + + @doc """ + GET /api/alerts -- Recent threshold-rule alerts emitted by + Hypatia.Watcher.Alerts (ring buffer, newest first). Powers the + dashboard alert ribbon and supports manual triage. + """ + get "/alerts" do + rows = + case Process.whereis(Hypatia.Watcher.Alerts) do + nil -> [] + _ -> Hypatia.Watcher.Alerts.recent() + end + + json(conn, 200, %{count: length(rows), rows: rows}) + end + + @doc """ + POST /api/alerts/ingest -- Federation ingress. Peer hypatia + instances POST their alerts here via the Peer sink. + + Auth: the auth_gate plug enforces a valid bearer token, so this + endpoint is only reachable when HYPATIA_API_BEARER_TOKEN is set + and the request carries it. Federation without shared auth is + refused at the gate, not here. + + Loop prevention: the ingested alert is tagged with + `metadata.federated_from = ` so the + Peer sink can skip it on broadcast and the dashboard can + attribute it. + """ + post "/alerts/ingest" do + {:ok, body, conn} = Plug.Conn.read_body(conn) + + case Jason.decode(body) do + {:ok, payload} when is_map(payload) -> + peer_id = + case Plug.Conn.get_req_header(conn, "x-forwarded-for") do + [host | _] -> host + _ -> conn.remote_ip |> :inet.ntoa() |> to_string() + end + + alert = + %{ + rule: parse_atom(payload["rule"]), + severity: parse_atom(payload["severity"]), + summary: payload["summary"] || "(no summary)", + metadata: payload["metadata"] || %{}, + at: payload["at"] || System.system_time(:millisecond) + } + + case Process.whereis(Hypatia.Watcher.Alerts) do + nil -> json(conn, 503, %{error: "alerts_unavailable"}) + _ -> Hypatia.Watcher.Alerts.ingest_federated(alert, peer_id) + end + + json(conn, 202, %{ok: true, peer: peer_id}) + + _ -> + json(conn, 400, %{error: "invalid_json"}) + end + end + + defp parse_atom(value) when is_atom(value), do: value + defp parse_atom(value) when is_binary(value) do + String.to_existing_atom(value) + rescue + ArgumentError -> :unknown + end + defp parse_atom(_), do: :unknown + + @doc """ + GET /api/events -- Server-Sent Events stream of telemetry as it + fires. Each event arrives as + + event: hypatia.scan.complete + data: {"measurements": {...}, "metadata": {...}, "at": ms} + + Optional `?events=hypatia.scan.complete,hypatia.outcome.recorded` + filter narrows the stream to specific event kinds. + + Heartbeats every 15s as comment lines (`: keepalive`) defeat proxy + idle-timeouts. The handler exits cleanly when the client disconnects + (Bandit closes the chunked response). + """ + get "/events" do + conn = Plug.Conn.fetch_query_params(conn) + filter = parse_event_filter(conn.query_params["events"]) + + case filter do + {:error, msg} -> + json(conn, 400, %{error: msg}) + + events_to_listen -> + Hypatia.Watcher.subscribe(events: events_to_listen) + + conn = + conn + |> put_resp_header("content-type", "text/event-stream") + |> put_resp_header("cache-control", "no-cache") + |> put_resp_header("connection", "keep-alive") + |> send_chunked(200) + + sse_loop(conn, schedule_heartbeat()) + end + end + + # ─── SSE internals ───────────────────────────────────────────────────── + + defp sse_loop(conn, heartbeat_ref) do + receive do + {:hypatia_event, event, measurements, metadata, ts} -> + chunk_body = + "event: #{Enum.join(event, ".")}\n" <> + "data: " <> + Jason.encode!(%{ + measurements: measurements, + metadata: jsonable_metadata(metadata), + at: ts + }) <> "\n\n" + + case Plug.Conn.chunk(conn, chunk_body) do + {:ok, conn} -> sse_loop(conn, heartbeat_ref) + {:error, _closed} -> conn + end + + :heartbeat -> + case Plug.Conn.chunk(conn, ": keepalive\n\n") do + {:ok, conn} -> sse_loop(conn, schedule_heartbeat()) + {:error, _closed} -> conn + end + end + end + + defp schedule_heartbeat do + Process.send_after(self(), :heartbeat, 15_000) + end + + # `:all` means subscribe to everything. A list of dotted-string + # event names means subscribe to that subset. Anything else is a + # client bug. + defp parse_event_filter(nil), do: :all + defp parse_event_filter(""), do: :all + + defp parse_event_filter(csv) do + csv + |> String.split(",", trim: true) + |> Enum.map(&parse_event_name/1) + |> Enum.reduce_while([], fn + {:ok, evt}, acc -> {:cont, [evt | acc]} + {:error, raw}, _acc -> {:halt, {:error, "unknown_event_name: #{raw}"}} + end) + |> case do + {:error, _} = err -> err + list -> Enum.reverse(list) + end + end + + defp parse_event_name(name) do + parts = name |> String.split(".") |> Enum.map(&safe_to_atom/1) + + if Enum.all?(parts, &is_atom/1) and parts in Hypatia.Telemetry.all_events() do + {:ok, parts} + else + {:error, name} + end + end + + defp safe_to_atom(s) do + String.to_existing_atom(s) + rescue + ArgumentError -> :__invalid__ + end + + # Metadata may contain pids/refs/funs that Jason can't encode. + # Coerce them to strings so the stream is always valid JSON. + defp jsonable_metadata(metadata) when is_map(metadata) do + Map.new(metadata, fn {k, v} -> {k, jsonable_value(v)} end) + end + + defp jsonable_value(v) + when is_binary(v) or is_number(v) or is_boolean(v) or is_atom(v) or is_nil(v), + do: v + + defp jsonable_value(v) when is_list(v), do: Enum.map(v, &jsonable_value/1) + defp jsonable_value(v) when is_map(v), do: jsonable_metadata(v) + defp jsonable_value(v), do: inspect(v) + match _ do json(conn, 404, %{error: "not_found"}) end # ─── Plug ────────────────────────────────────────────────────────────── + # ─── Auth gate ───────────────────────────────────────────────────────── + # + # If HYPATIA_API_BEARER_TOKEN is set, any /api/* request must carry a + # matching Authorization: Bearer header. The token + loopback + # checks compose: with neither, /api is loopback-only. With both, /api + # is openable to non-local callers provided they present the token. + # + # Token comparison uses Plug.Crypto.secure_compare/2 so timing attacks + # can't enumerate the secret. + defp auth_gate(conn, _opts) do + case System.get_env("HYPATIA_API_BEARER_TOKEN") do + nil -> + conn + + "" -> + conn + + expected -> + case get_bearer(conn) do + {:ok, presented} -> + if secure_compare(expected, presented) do + # If the token is valid, the caller is implicitly allowed + # past the loopback gate too — explicit authentication is + # at least as strong as IP-based filtering. + Plug.Conn.put_private(conn, :hypatia_auth_passed, true) + else + conn + |> put_resp_content_type("application/json") + |> send_resp( + 401, + Jason.encode!(%{error: "invalid_token", hint: "Authorization: Bearer "}) + ) + |> halt() + end + + :error -> + conn + |> put_resp_content_type("application/json") + |> send_resp( + 401, + Jason.encode!(%{ + error: "missing_token", + hint: + "Hypatia API requires Authorization: Bearer when " <> + "HYPATIA_API_BEARER_TOKEN is set on the server." + }) + ) + |> halt() + end + end + end + + defp get_bearer(conn) do + case Plug.Conn.get_req_header(conn, "authorization") do + ["Bearer " <> token] -> {:ok, token} + ["bearer " <> token] -> {:ok, token} + _ -> :error + end + end + + # Constant-time string comparison. Plug.Crypto provides this in the + # production dep set; fall back to a hand-rolled equivalent for + # environments where Plug.Crypto isn't loaded (e.g. escript builds). + defp secure_compare(a, b) when is_binary(a) and is_binary(b) do + if Code.ensure_loaded?(Plug.Crypto) and function_exported?(Plug.Crypto, :secure_compare, 2) do + Plug.Crypto.secure_compare(a, b) + else + byte_size(a) == byte_size(b) and + a |> :binary.bin_to_list() |> Enum.zip(:binary.bin_to_list(b)) |> Enum.reduce(0, fn {x, y}, acc -> acc ||| Bitwise.bxor(x, y) end) == 0 + end + end + defp loopback_only(conn, _opts) do + if conn.private[:hypatia_auth_passed] do + # Caller authenticated via bearer; loopback check is unnecessary. + conn + else + legacy_loopback_only(conn, []) + end + end + + defp legacy_loopback_only(conn, _opts) do cond do System.get_env("HYPATIA_API_ALLOW_NONLOCAL") == "true" -> Logger.warning( diff --git a/lib/hypatia/web/dashboard.ex b/lib/hypatia/web/dashboard.ex new file mode 100644 index 00000000..a79eb46c --- /dev/null +++ b/lib/hypatia/web/dashboard.ex @@ -0,0 +1,294 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Web.Dashboard do + @moduledoc """ + Single-page operational dashboard for Hypatia. + + Served at `GET /` by `Hypatia.Web.Router`. Plain HTML + vanilla JS + + CSS — no build step, no framework, no node_modules. The page + polls `/api/status` every 2s for the counters / queue depths, and + opens an `EventSource` connection to `/api/events` for live + telemetry as it fires. + + Loopback-only by design — uses the same gate as the rest of /api/* + (the gate plug runs in `ApiRouter`; the dashboard makes XHR/SSE + calls to those endpoints, so a non-local browser would be rejected + by the API endpoints even if it reached the dashboard). + + The HTML is generated at compile-time (module attribute) so there's + no template-rendering cost on each request and no on-disk asset + file to misplace. + """ + + import Plug.Conn + + @html ~S""" + + + + + Hypatia — Live Watcher + + + + +
+
+

Hypatia Watcher

+
+ Uptime — + · Dropped 0 + · stream offline +
+
+
never updated
+
+ +
+
+

Events / 5 minutes

+
loading…
+
+
+

Events / 1 hour

+
loading…
+
+
+

GenServer queue depths

+
loading…
+
+
+

Recipe health (actionable rows first)

+
loading…
+
+
+ +
+

Live telemetry stream

+
connecting…
+
+ + + + + """ + + @doc """ + Plug entry point for the dashboard HTML. + """ + def call(conn, _opts) do + conn + |> put_resp_content_type("text/html; charset=utf-8") + |> put_resp_header("cache-control", "no-store") + |> send_resp(200, @html) + end + + def init(opts), do: opts +end diff --git a/lib/hypatia/web/metrics.ex b/lib/hypatia/web/metrics.ex new file mode 100644 index 00000000..7ea09ae9 --- /dev/null +++ b/lib/hypatia/web/metrics.ex @@ -0,0 +1,211 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Web.Metrics do + @moduledoc """ + Prometheus text-format `/metrics` exposition. + + Hand-rolled rather than via `:telemetry_metrics_prometheus` so we + don't add a runtime dependency. The text format is simple and + stable — `# HELP` + `# TYPE` + value rows. + + Surface: + + hypatia_events_total{event="hypatia.scan.complete", window="5m"} N + hypatia_events_total{event="hypatia.dispatch.decision", window="1h"} N + ... + hypatia_queue_depth{process="Hypatia.Watcher"} N + hypatia_watcher_dropped_total N + hypatia_watcher_uptime_seconds N + hypatia_recipe_verification_rate{recipe="recipe-foo"} 0.93 + hypatia_recipe_dispatches_total{recipe="recipe-foo"} 1234 + hypatia_recipe_quarantine_candidates N + hypatia_recipe_degraded N + + Bound publicly (NOT loopback-only) because Prometheus scrapers + routinely run on a different host than the app — there's no + operational data in the metric body that isn't already implied by + the dashboard. If this assumption changes, move under /api. + + ## Prometheus scrape config + + scrape_configs: + - job_name: hypatia + metrics_path: /metrics + static_configs: + - targets: ['hypatia.internal:9090'] + """ + + import Plug.Conn + + def call(conn, _opts) do + body = render() + + conn + |> put_resp_content_type("text/plain; version=0.0.4; charset=utf-8") + |> send_resp(200, body) + end + + def init(opts), do: opts + + @doc """ + Render the Prometheus text-format exposition. Public so the unit + test can assert against the body without going through Plug. + """ + def render do + [ + render_events_total(), + render_queue_depths(), + render_watcher_meta(), + render_recipe_health() + ] + |> Enum.join("\n") + end + + # ─── Event counters per window ───────────────────────────────────────── + + defp render_events_total do + lines = + for window <- [:m5, :h1, :d1], + {event, count} <- safe_counts(window) do + ~s|hypatia_events_total{event="#{Enum.join(event, ".")}",window="#{window_label(window)}"} #{count}| + end + + [ + "# HELP hypatia_events_total Telemetry events counted in the named rolling window", + "# TYPE hypatia_events_total gauge" + | lines + ] + |> Enum.join("\n") + end + + defp window_label(:m5), do: "5m" + defp window_label(:h1), do: "1h" + defp window_label(:d1), do: "1d" + + defp safe_counts(window) do + Hypatia.Watcher.counts(window) + rescue + _ -> %{} + catch + _, _ -> %{} + end + + # ─── GenServer queue depths ──────────────────────────────────────────── + + defp render_queue_depths do + depths = safe_queue_depths() + + lines = + for {name, depth} <- depths, is_integer(depth) do + # Strip module-inspect formatting (e.g. "Elixir.Hypatia.Watcher" + # → "Hypatia.Watcher") so Prometheus labels are tidy. + clean = name |> to_string() |> String.replace_leading("Elixir.", "") + ~s|hypatia_queue_depth{process="#{clean}"} #{depth}| + end + + [ + "# HELP hypatia_queue_depth GenServer mailbox depth for supervised processes", + "# TYPE hypatia_queue_depth gauge" + | lines + ] + |> Enum.join("\n") + end + + defp safe_queue_depths do + Hypatia.Watcher.queue_depths() + rescue + _ -> %{} + catch + _, _ -> %{} + end + + # ─── Watcher meta (dropped events, uptime) ───────────────────────────── + + defp render_watcher_meta do + snap = safe_snapshot() + dropped = Map.get(snap, :dropped_events, 0) + uptime = Map.get(snap, :uptime_seconds, 0) + + """ + # HELP hypatia_watcher_dropped_total Telemetry events dropped under back-pressure + # TYPE hypatia_watcher_dropped_total counter + hypatia_watcher_dropped_total #{dropped} + # HELP hypatia_watcher_uptime_seconds Seconds since the Watcher started + # TYPE hypatia_watcher_uptime_seconds gauge + hypatia_watcher_uptime_seconds #{uptime} + """ + |> String.trim_trailing() + end + + defp safe_snapshot do + Hypatia.Watcher.snapshot() + rescue + _ -> %{} + catch + _, _ -> %{} + end + + # ─── Recipe health ───────────────────────────────────────────────────── + + defp render_recipe_health do + rows = safe_recipe_health() + + rate_lines = + for r <- rows, is_float(r.verification.rate) do + ~s|hypatia_recipe_verification_rate{recipe="#{escape(r.recipe_id)}"} #{r.verification.rate}| + end + + dispatch_lines = + for r <- rows do + ~s|hypatia_recipe_dispatches_total{recipe="#{escape(r.recipe_id)}"} #{r.dispatches}| + end + + qc = Enum.count(rows, &(&1.status == :quarantine_candidate)) + deg = Enum.count(rows, &(&1.status == :degraded)) + + sections = [ + [ + "# HELP hypatia_recipe_verification_rate Fraction of recipe successes confirmed clean by post-fix re-scan", + "# TYPE hypatia_recipe_verification_rate gauge" + | rate_lines + ], + [ + "# HELP hypatia_recipe_dispatches_total Total dispatch attempts recorded for the recipe", + "# TYPE hypatia_recipe_dispatches_total counter" + | dispatch_lines + ], + [ + "# HELP hypatia_recipe_quarantine_candidates Recipes auto-quarantined on verification-rate gate", + "# TYPE hypatia_recipe_quarantine_candidates gauge", + "hypatia_recipe_quarantine_candidates #{qc}", + "# HELP hypatia_recipe_degraded Recipes below the degraded threshold (but above quarantine)", + "# TYPE hypatia_recipe_degraded gauge", + "hypatia_recipe_degraded #{deg}" + ] + ] + + sections + |> Enum.map(&Enum.join(&1, "\n")) + |> Enum.join("\n") + end + + defp safe_recipe_health do + Hypatia.OutcomeTracker.recipe_health() + rescue + _ -> [] + catch + _, _ -> [] + end + + # Prometheus label values must escape backslash, double-quote, and + # newline per the exposition spec. + defp escape(value) when is_binary(value) do + value + |> String.replace("\\", "\\\\") + |> String.replace("\"", "\\\"") + |> String.replace("\n", "\\n") + end + + defp escape(value), do: value |> to_string() |> escape() +end diff --git a/lib/hypatia/web/router.ex b/lib/hypatia/web/router.ex index 9084f4dc..7893a670 100644 --- a/lib/hypatia/web/router.ex +++ b/lib/hypatia/web/router.ex @@ -28,6 +28,17 @@ defmodule Hypatia.Web.Router do plug :match plug :dispatch + @doc """ + GET / -- Single-page live operational dashboard. HTML + vanilla JS, + polls /api/status and EventSource-streams /api/events. The dashboard + itself is publicly reachable; the data endpoints it calls are + loopback-only (gated in ApiRouter), so a non-local browser would + render the chrome but get 403 from the XHR/SSE calls. + """ + get "/" do + Hypatia.Web.Dashboard.call(conn, []) + end + @doc """ GET /health -- Basic health check for the HTTP endpoint. """ @@ -43,6 +54,16 @@ defmodule Hypatia.Web.Router do |> send_resp(200, Jason.encode!(health)) end + @doc """ + GET /metrics -- Prometheus text-format exposition. Publicly + reachable (NOT loopback-only) because scrapers routinely run on a + different host; there's no operational data in the metric body + that isn't already implied by the dashboard's existence. + """ + get "/metrics" do + Hypatia.Web.Metrics.call(conn, []) + end + # /api/* is gated to loopback in Hypatia.Web.ApiRouter — keeps # operational data off the public surface while leaving /health # reachable for container orchestrators. diff --git a/test/alerts_test.exs b/test/alerts_test.exs new file mode 100644 index 00000000..15603dc3 --- /dev/null +++ b/test/alerts_test.exs @@ -0,0 +1,117 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.AlertsTest do + # async: false — the Alerts GenServer is a named singleton that + # registers global handlers; concurrent tests would observe each + # other's emissions. + use ExUnit.Case, async: false + + alias Hypatia.Watcher.Alerts + alias Hypatia.Telemetry, as: T + + setup do + case Process.whereis(Hypatia.Watcher.PubSub) do + nil -> + {:ok, pid} = Registry.start_link(keys: :duplicate, name: Hypatia.Watcher.PubSub) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(Hypatia.Watcher) do + nil -> + {:ok, pid} = Hypatia.Watcher.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(Hypatia.Watcher.Alerts) do + nil -> + {:ok, pid} = Alerts.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + # No external sinks by default — test isolation. + System.delete_env("HYPATIA_ALERT_WEBHOOK_URL") + System.delete_env("HYPATIA_ALERT_LOG_FILE") + + :ok + end + + describe "quarantine_triggered rule" do + test "emits an alert on the corresponding telemetry event" do + T.quarantine_triggered(kind: :recipe, id: "bad-recipe-1", reason: "rate", level: :auto) + # The cast to Alerts is async; wait a beat. + Process.sleep(100) + + alerts = Alerts.recent() + [head | _] = alerts + + assert head.rule == :quarantine_triggered + assert head.severity == :high + assert head.summary =~ "bad-recipe-1" + end + + test "deduplicates within the dedup window" do + T.quarantine_triggered(kind: :recipe, id: "dup-recipe", reason: "rate", level: :auto) + Process.sleep(50) + T.quarantine_triggered(kind: :recipe, id: "dup-recipe", reason: "rate", level: :auto) + Process.sleep(50) + T.quarantine_triggered(kind: :recipe, id: "dup-recipe", reason: "rate", level: :auto) + Process.sleep(50) + + matching = + Alerts.recent() + |> Enum.filter(fn a -> a.rule == :quarantine_triggered and a.summary =~ "dup-recipe" end) + + assert length(matching) == 1 + end + + test "different ids are not deduped against each other" do + T.quarantine_triggered(kind: :recipe, id: "id-a-#{System.unique_integer()}", reason: "r", level: :auto) + T.quarantine_triggered(kind: :recipe, id: "id-b-#{System.unique_integer()}", reason: "r", level: :auto) + Process.sleep(100) + + recent_summaries = Alerts.recent() |> Enum.map(& &1.summary) + a_count = Enum.count(recent_summaries, &(&1 =~ "id-a-")) + b_count = Enum.count(recent_summaries, &(&1 =~ "id-b-")) + assert a_count >= 1 + assert b_count >= 1 + end + end + + describe "soundness_violation rule" do + test "fires immediately, with severity :critical" do + T.soundness_violation( + rule_module: "code_safety", + rule_id: "elixir_system_shell", + fixture: "test/soundness/fixtures/code_safety/elixir_system_shell.ex" + ) + + Process.sleep(100) + + [head | _] = Alerts.recent() + assert head.rule == :soundness_violation + assert head.severity == :critical + assert head.summary =~ "elixir_system_shell" + end + end + + describe "events_dropped rule (tick path)" do + test "fires when watcher dropped_events delta > 0" do + # We can't directly mutate the watcher's internal state, but we can + # craft the rule's input by exercising the tick path against a + # snapshot the watcher returns. Easiest: assert the tick is callable + # and recent() is queryable without crashing. + assert :ok = Alerts.tick_now() + assert is_list(Alerts.recent()) + end + end +end diff --git a/test/anomaly_detector_test.exs b/test/anomaly_detector_test.exs new file mode 100644 index 00000000..de295f3e --- /dev/null +++ b/test/anomaly_detector_test.exs @@ -0,0 +1,124 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.AnomalyDetectorTest do + use ExUnit.Case, async: false + + alias Hypatia.Watcher.AnomalyDetector + alias Hypatia.Telemetry, as: T + + setup do + case Process.whereis(Hypatia.Watcher.PubSub) do + nil -> + {:ok, pid} = Registry.start_link(keys: :duplicate, name: Hypatia.Watcher.PubSub) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(Hypatia.Watcher) do + nil -> + {:ok, pid} = Hypatia.Watcher.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(AnomalyDetector) do + nil -> + {:ok, pid} = AnomalyDetector.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + :ok + end + + describe "outcome ingestion" do + test "rolling window accumulates outcome.recorded events" do + for _ <- 1..5 do + T.outcome_recorded(recipe_id: "r", repo: "x", outcome: "success", verification: "verified") + end + + Process.sleep(100) + + window = AnomalyDetector.window() + assert length(window) >= 5 + end + end + + describe "evaluate (insufficient data)" do + test "below min_outcomes_for_alert, tick is a no-op" do + # Fresh detector — fewer than 30 outcomes + assert :ok = AnomalyDetector.tick_now() + end + end + + describe "evaluate (clear regression)" do + test "emits hypatia.anomaly.detected when recent rate drops 100% → 0%" do + :ok = :telemetry.attach("anomaly-test-handler", + [:hypatia, :anomaly, :detected], + fn _event, measurements, metadata, _config -> + send(self(), {:caught, measurements, metadata}) + end, + nil) + + on_exit(fn -> :telemetry.detach("anomaly-test-handler") end) + + # 170 baseline successes + 30 recent failures = clear regression + for _ <- 1..170 do + T.outcome_recorded(recipe_id: "r", repo: "x", outcome: "success", verification: "verified") + end + + for _ <- 1..30 do + T.outcome_recorded(recipe_id: "r", repo: "x", outcome: "failure", verification: "unverified") + end + + # Let the cast queue drain + Process.sleep(200) + + :ok = AnomalyDetector.tick_now() + + # Handler runs in the *emitting* process (the detector). It sends + # to whatever pid is captured at attach time -- which was the + # test process. So we receive in this test's mailbox. + assert_receive {:caught, measurements, metadata}, 1000 + + assert measurements.recent_rate < 0.5 + assert measurements.baseline_rate > 0.5 + assert measurements.sigma_distance > 2.0 + assert metadata.kind == :success_rate_drop + end + end + + describe "evaluate (stable history)" do + test "does NOT emit when recent rate matches baseline" do + caller = self() + + :ok = :telemetry.attach("anomaly-stable-handler", + [:hypatia, :anomaly, :detected], + fn _event, _measurements, _metadata, _config -> + send(caller, :unexpected_anomaly) + end, + nil) + + on_exit(fn -> :telemetry.detach("anomaly-stable-handler") end) + + # Mixed but stable history: 80% success across both halves. + for i <- 1..200 do + outcome = if rem(i, 5) == 0, do: "failure", else: "success" + T.outcome_recorded(recipe_id: "r", repo: "x", outcome: outcome, verification: "verified") + end + + Process.sleep(200) + + :ok = AnomalyDetector.tick_now() + + refute_receive :unexpected_anomaly, 200 + end + end +end diff --git a/test/api_router_test.exs b/test/api_router_test.exs index 5d1e9c22..1ce9fe6a 100644 --- a/test/api_router_test.exs +++ b/test/api_router_test.exs @@ -22,6 +22,76 @@ defmodule Hypatia.Web.ApiRouterTest do :ok end + describe "bearer auth gate" do + test "no token configured → no auth required (legacy behaviour)" do + System.delete_env("HYPATIA_API_BEARER_TOKEN") + conn = build_conn(:get, "/status", {127, 0, 0, 1}) + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 200 + end + + test "token configured + valid Bearer header → 200" do + System.put_env("HYPATIA_API_BEARER_TOKEN", "test-secret-abc123") + on_exit(fn -> System.delete_env("HYPATIA_API_BEARER_TOKEN") end) + + conn = + build_conn(:get, "/status", {127, 0, 0, 1}) + |> Plug.Conn.put_req_header("authorization", "Bearer test-secret-abc123") + + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 200 + end + + test "token configured + wrong Bearer header → 401" do + System.put_env("HYPATIA_API_BEARER_TOKEN", "test-secret-abc123") + on_exit(fn -> System.delete_env("HYPATIA_API_BEARER_TOKEN") end) + + conn = + build_conn(:get, "/status", {127, 0, 0, 1}) + |> Plug.Conn.put_req_header("authorization", "Bearer wrong-token") + + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 401 + body = Jason.decode!(conn.resp_body) + assert body["error"] == "invalid_token" + end + + test "token configured + missing header → 401 missing_token" do + System.put_env("HYPATIA_API_BEARER_TOKEN", "test-secret-abc123") + on_exit(fn -> System.delete_env("HYPATIA_API_BEARER_TOKEN") end) + + conn = build_conn(:get, "/status", {127, 0, 0, 1}) + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 401 + body = Jason.decode!(conn.resp_body) + assert body["error"] == "missing_token" + end + + test "valid token from non-loopback IP → 200 (bearer beats loopback)" do + System.put_env("HYPATIA_API_BEARER_TOKEN", "test-secret-abc123") + on_exit(fn -> System.delete_env("HYPATIA_API_BEARER_TOKEN") end) + + conn = + build_conn(:get, "/status", {10, 1, 2, 3}) + |> Plug.Conn.put_req_header("authorization", "Bearer test-secret-abc123") + + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 200 + end + + test "lowercase 'bearer' prefix is accepted" do + System.put_env("HYPATIA_API_BEARER_TOKEN", "test-secret-abc123") + on_exit(fn -> System.delete_env("HYPATIA_API_BEARER_TOKEN") end) + + conn = + build_conn(:get, "/status", {127, 0, 0, 1}) + |> Plug.Conn.put_req_header("authorization", "bearer test-secret-abc123") + + conn = ApiRouter.call(conn, ApiRouter.init([])) + assert conn.status == 200 + end + end + describe "loopback gate" do test "127.0.0.1 caller is allowed through" do conn = build_conn(:get, "/status", {127, 0, 0, 1}) @@ -124,6 +194,31 @@ defmodule Hypatia.Web.ApiRouterTest do end end + describe "GET /recipes/:id" do + test "returns 404 for an unknown recipe id" do + conn = build_conn(:get, "/recipes/does-not-exist-recipe-id", {127, 0, 0, 1}) + conn = ApiRouter.call(conn, ApiRouter.init([])) + + assert conn.status == 404 + body = Jason.decode!(conn.resp_body) + assert body["error"] == "recipe_not_found" + end + end + + describe "GET /quarantine" do + test "returns the recipe + bot quarantine roster" do + conn = build_conn(:get, "/quarantine", {127, 0, 0, 1}) + conn = ApiRouter.call(conn, ApiRouter.init([])) + + assert conn.status == 200 + body = Jason.decode!(conn.resp_body) + assert Map.has_key?(body, "recipes") + assert Map.has_key?(body, "bots") + assert is_list(body["recipes"]["rows"]) + assert is_map(body["bots"]["entries"]) + end + end + defp build_conn(method, path, remote_ip) do conn(method, path) |> Map.put(:remote_ip, remote_ip) diff --git a/test/dashboard_test.exs b/test/dashboard_test.exs new file mode 100644 index 00000000..277b132a --- /dev/null +++ b/test/dashboard_test.exs @@ -0,0 +1,53 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Web.DashboardTest do + use ExUnit.Case, async: true + use Plug.Test + + alias Hypatia.Web.Dashboard + + describe "GET /" do + test "serves the dashboard HTML" do + conn = conn(:get, "/") |> Dashboard.call([]) + + assert conn.status == 200 + + ctype = + conn + |> get_resp_header("content-type") + |> List.first() + + assert ctype =~ "text/html" + assert conn.resp_body =~ "Hypatia Watcher" + assert conn.resp_body =~ "EventSource" + assert conn.resp_body =~ "/api/status" + assert conn.resp_body =~ "/api/events" + end + + test "sets no-store cache header so refreshes always get fresh dashboard" do + conn = conn(:get, "/") |> Dashboard.call([]) + + cache_control = + conn + |> get_resp_header("cache-control") + |> List.first() + + assert cache_control == "no-store" + end + + test "references every known telemetry event name in the JS event handler list" do + conn = conn(:get, "/") |> Dashboard.call([]) + + for event <- Hypatia.Telemetry.all_events() do + dotted = Enum.join(event, ".") + + assert conn.resp_body =~ dotted, + "dashboard JS handler list missing event '#{dotted}' — " <> + "if you add a new telemetry event in Hypatia.Telemetry, " <> + "add it to the `known` array in dashboard.ex so the " <> + "stream actually renders it" + end + end + end +end diff --git a/test/metrics_test.exs b/test/metrics_test.exs new file mode 100644 index 00000000..8aa08e63 --- /dev/null +++ b/test/metrics_test.exs @@ -0,0 +1,113 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Web.MetricsTest do + use ExUnit.Case, async: false + use Plug.Test + + alias Hypatia.Web.Metrics + alias Hypatia.Telemetry, as: T + + setup do + case Process.whereis(Hypatia.Watcher) do + nil -> + {:ok, pid} = Hypatia.Watcher.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + :ok + end + + describe "render/0 — exposition format" do + test "emits HELP + TYPE lines for each metric family" do + output = Metrics.render() + + assert output =~ "# HELP hypatia_events_total " + assert output =~ "# TYPE hypatia_events_total gauge" + assert output =~ "# HELP hypatia_queue_depth " + assert output =~ "# TYPE hypatia_queue_depth gauge" + assert output =~ "# HELP hypatia_watcher_dropped_total " + assert output =~ "# TYPE hypatia_watcher_dropped_total counter" + assert output =~ "# HELP hypatia_watcher_uptime_seconds " + assert output =~ "# TYPE hypatia_watcher_uptime_seconds gauge" + assert output =~ "# HELP hypatia_recipe_verification_rate " + assert output =~ "# HELP hypatia_recipe_quarantine_candidates " + assert output =~ "# HELP hypatia_recipe_degraded " + end + + test "event counter labels include window and event" do + T.scan_complete(50, 3, path: "/tmp/x", severity_floor: "low") + Process.sleep(50) + + output = Metrics.render() + + # Should appear in all three windows: 5m / 1h / 1d + assert output =~ + ~r/hypatia_events_total\{event="hypatia.scan.complete",window="5m"\} \d+/ + + assert output =~ + ~r/hypatia_events_total\{event="hypatia.scan.complete",window="1h"\} \d+/ + + assert output =~ + ~r/hypatia_events_total\{event="hypatia.scan.complete",window="1d"\} \d+/ + end + + test "queue depth label strips Elixir. prefix" do + # If no supervisor is running there are no queue lines, but if there + # is one the format must be tidy. Don't require content — assert + # negation: NO line contains "Elixir." in a label. + output = Metrics.render() + refute output =~ ~r/process="Elixir\./ + end + + test "label values are escaped per the Prometheus spec" do + # Backslashes and quotes are the dangerous chars; we don't expect + # recipe ids to contain them, but assert by direct test of the + # output format on something the renderer emits. + output = Metrics.render() + lines = String.split(output, "\n") + + # Any line with a label must close all its quote pairs. + Enum.each(lines, fn line -> + if line =~ "{" do + # Count unescaped quotes within braces only. + inside = + line + |> String.split("{", parts: 2) + |> List.last() + |> String.split("}", parts: 2) + |> List.first() + + # Even number of unescaped " in label section. + unescaped_quotes = + inside + |> String.replace("\\\"", "") + |> String.graphemes() + |> Enum.count(&(&1 == "\"")) + + assert rem(unescaped_quotes, 2) == 0, + "unbalanced quotes in label section: #{line}" + end + end) + end + end + + describe "GET /metrics" do + test "responds with prometheus text format content-type" do + conn = conn(:get, "/metrics") |> Metrics.call([]) + assert conn.status == 200 + + ctype = conn |> get_resp_header("content-type") |> List.first() + assert ctype =~ "text/plain" + assert ctype =~ "version=0.0.4" + end + + test "body is the same as Metrics.render/0" do + conn = conn(:get, "/metrics") |> Metrics.call([]) + assert conn.resp_body == Metrics.render() + end + end +end diff --git a/test/persistence_test.exs b/test/persistence_test.exs new file mode 100644 index 00000000..be3d4ed3 --- /dev/null +++ b/test/persistence_test.exs @@ -0,0 +1,115 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.Watcher.PersistenceTest do + use ExUnit.Case, async: false + + alias Hypatia.Watcher.Persistence + alias Hypatia.Telemetry, as: T + + setup do + # Isolate from the on-disk verisim metrics by pointing the module + # at a per-test tmp dir. + tmp = Path.join(System.tmp_dir!(), "hypatia-persist-#{System.unique_integer([:positive])}") + File.mkdir_p!(tmp) + Application.put_env(:hypatia, :metrics_path, tmp) + + case Process.whereis(Hypatia.Watcher.PubSub) do + nil -> + {:ok, pid} = Registry.start_link(keys: :duplicate, name: Hypatia.Watcher.PubSub) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(Hypatia.Watcher) do + nil -> + {:ok, pid} = Hypatia.Watcher.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + case Process.whereis(Persistence) do + nil -> + {:ok, pid} = Persistence.start_link([]) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + on_exit(fn -> + Application.delete_env(:hypatia, :metrics_path) + File.rm_rf!(tmp) + end) + + {:ok, tmp_dir: tmp} + end + + describe "snapshot_now/0" do + test "writes a JSONL record to today's file", %{tmp_dir: tmp} do + T.scan_complete(50, 3, path: "/tmp/x", severity_floor: "low") + Process.sleep(50) + + {:ok, path, record} = Persistence.snapshot_now() + + assert String.starts_with?(path, tmp) + assert String.ends_with?(path, ".jsonl") + assert File.exists?(path) + + [line | _] = File.read!(path) |> String.split("\n", trim: true) |> Enum.reverse() + decoded = Jason.decode!(line) + + assert decoded["at_ms"] == record.at_ms + assert is_map(decoded["counts_m5"]) + assert is_map(decoded["counts_h1"]) + assert is_map(decoded["counts_d1"]) + assert is_map(decoded["queue_depths"]) + assert is_map(decoded["recipe_health_summary"]) + assert is_integer(decoded["alert_count"]) + end + + test "multiple snapshots produce multiple JSONL lines" do + Persistence.snapshot_now() + Persistence.snapshot_now() + Persistence.snapshot_now() + + path = Persistence.current_file() + lines = File.read!(path) |> String.split("\n", trim: true) + + assert length(lines) >= 3 + + Enum.each(lines, fn line -> + assert {:ok, _} = Jason.decode(line), + "Each persistence line must be valid JSON: #{line}" + end) + end + + test "counts_m5 keys are dotted-string event names, not lists" do + T.dispatch_decision(0.9, strategy: :review, tier: :eliminate, recipe_id: "r", repo: "x") + Process.sleep(50) + + {:ok, _path, record} = Persistence.snapshot_now() + + keys = Map.keys(record.counts_m5) + + Enum.each(keys, fn k -> + assert is_binary(k), + "counts_m5 keys must be strings (got #{inspect(k)})" + + assert k =~ ".", + "counts_m5 keys must be dotted (got #{inspect(k)})" + end) + end + end + + describe "current_file/0" do + test "today's path uses YYYY-MM-DD.jsonl" do + path = Persistence.current_file() + assert path =~ ~r/\d{4}-\d{2}-\d{2}\.jsonl$/ + end + end +end diff --git a/test/sarif_test.exs b/test/sarif_test.exs new file mode 100644 index 00000000..3f265f1b --- /dev/null +++ b/test/sarif_test.exs @@ -0,0 +1,166 @@ +# SPDX-License-Identifier: MPL-2.0 +# Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) + +defmodule Hypatia.SARIFTest do + use ExUnit.Case, async: true + + alias Hypatia.SARIF + + describe "from_findings/2 — shape" do + test "empty findings produce a valid SARIF document with an empty results array" do + doc = SARIF.from_findings([], "/tmp") + + assert doc["$schema"] =~ "sarif-2.1.0" + assert doc["version"] == "2.1.0" + [run] = doc["runs"] + assert run["tool"]["driver"]["name"] == "Hypatia" + assert run["results"] == [] + assert run["tool"]["driver"]["rules"] == [] + end + + test "one finding produces one result + one rule" do + doc = + SARIF.from_findings( + [ + %{ + severity: "critical", + rule_module: "code_safety", + type: "elixir_system_shell", + file: "lib/foo.ex", + reason: "System.shell with interpolation -- shell injection risk" + } + ], + "/tmp/repo" + ) + + [run] = doc["runs"] + [result] = run["results"] + [rule] = run["tool"]["driver"]["rules"] + + assert result["ruleId"] == "hypatia/code_safety/elixir_system_shell" + assert result["level"] == "error" + assert result["message"]["text"] =~ "shell injection" + assert result["locations"] |> List.first() |> get_in(["physicalLocation", "artifactLocation", "uri"]) == "lib/foo.ex" + assert is_binary(result["partialFingerprints"]["hypatiaFindingHash/v1"]) + + assert rule["id"] == "hypatia/code_safety/elixir_system_shell" + assert rule["defaultConfiguration"]["level"] == "error" + end + + test "dedups rules across multiple results with the same ruleId" do + findings = + for n <- 1..5 do + %{ + severity: "high", + rule_module: "code_safety", + type: "unwrap_without_check", + file: "lib/foo_#{n}.rs", + reason: "Rust unwrap" + } + end + + doc = SARIF.from_findings(findings, "/tmp") + [run] = doc["runs"] + + assert length(run["results"]) == 5 + assert length(run["tool"]["driver"]["rules"]) == 1 + end + + test "severity → level mapping matches the JS converter" do + cases = [ + {"critical", "error"}, + {"high", "error"}, + {"medium", "warning"}, + {"low", "note"}, + {"info", "note"}, + {"", "note"} + ] + + Enum.each(cases, fn {severity, expected_level} -> + finding = %{ + severity: severity, + rule_module: "x", + type: "y", + file: "a.ex", + reason: "r" + } + + [result] = SARIF.from_findings([finding], "/tmp") |> get_in(["runs", Access.at(0), "results"]) + assert result["level"] == expected_level, "severity #{inspect(severity)} → #{expected_level}" + end) + end + end + + describe "URI relativisation" do + test "absolute paths under the repo root become relative" do + doc = + SARIF.from_findings( + [ + %{severity: "high", rule_module: "x", type: "y", file: "/home/user/hypatia/lib/foo.ex", reason: "r"} + ], + "/home/user/hypatia" + ) + + uri = doc |> get_in(["runs", Access.at(0), "results", Access.at(0), "locations", Access.at(0), "physicalLocation", "artifactLocation", "uri"]) + assert uri == "lib/foo.ex" + end + + test "absolute paths OUTSIDE the repo root degrade to basename" do + doc = + SARIF.from_findings( + [%{severity: "high", rule_module: "x", type: "y", file: "/etc/passwd", reason: "r"}], + "/home/user/hypatia" + ) + + uri = doc |> get_in(["runs", Access.at(0), "results", Access.at(0), "locations", Access.at(0), "physicalLocation", "artifactLocation", "uri"]) + assert uri == "passwd" + end + + test "empty file becomes '.' (repo-level finding)" do + doc = + SARIF.from_findings( + [%{severity: "high", rule_module: "x", type: "y", file: "", reason: "r"}], + "/tmp" + ) + + uri = doc |> get_in(["runs", Access.at(0), "results", Access.at(0), "locations", Access.at(0), "physicalLocation", "artifactLocation", "uri"]) + assert uri == "." + end + end + + describe "render/2" do + test "produces parseable JSON" do + json = + SARIF.render( + [%{severity: "high", rule_module: "x", type: "y", file: "a.ex", reason: "r"}], + "/tmp" + ) + + assert {:ok, decoded} = Jason.decode(json) + assert decoded["version"] == "2.1.0" + end + end + + describe "fingerprints" do + test "same {ruleId, uri, type, reason} produces same fingerprint" do + f = %{severity: "high", rule_module: "x", type: "y", file: "a.ex", reason: "r"} + + [r1] = SARIF.from_findings([f], "/tmp") |> get_in(["runs", Access.at(0), "results"]) + [r2] = SARIF.from_findings([f], "/tmp") |> get_in(["runs", Access.at(0), "results"]) + + assert r1["partialFingerprints"]["hypatiaFindingHash/v1"] == + r2["partialFingerprints"]["hypatiaFindingHash/v1"] + end + + test "differing reason produces different fingerprint" do + f1 = %{severity: "high", rule_module: "x", type: "y", file: "a.ex", reason: "r1"} + f2 = %{severity: "high", rule_module: "x", type: "y", file: "a.ex", reason: "r2"} + + [r1] = SARIF.from_findings([f1], "/tmp") |> get_in(["runs", Access.at(0), "results"]) + [r2] = SARIF.from_findings([f2], "/tmp") |> get_in(["runs", Access.at(0), "results"]) + + refute r1["partialFingerprints"]["hypatiaFindingHash/v1"] == + r2["partialFingerprints"]["hypatiaFindingHash/v1"] + end + end +end diff --git a/test/watcher_test.exs b/test/watcher_test.exs index 43952d2d..c2dfd324 100644 --- a/test/watcher_test.exs +++ b/test/watcher_test.exs @@ -106,4 +106,92 @@ defmodule Hypatia.WatcherTest do assert is_map(depths) end end + + describe "persistence across restart" do + test "restart restores ETS counters + recent-event tail" do + # Use a per-test tmp dir for isolation. + tmp = Path.join(System.tmp_dir!(), "watcher-restart-#{System.unique_integer([:positive])}") + File.mkdir_p!(tmp) + Application.put_env(:hypatia, :watcher_persist_path, tmp) + + on_exit(fn -> + Application.delete_env(:hypatia, :watcher_persist_path) + File.rm_rf!(tmp) + end) + + # Bring down any existing watcher so we can boot fresh into the tmp. + if pid = Process.whereis(Watcher) do + if Process.alive?(pid), do: GenServer.stop(pid) + end + + {:ok, pid1} = Watcher.start_link([]) + T.scan_complete(99, 3, path: "/tmp/x", severity_floor: "low") + Process.sleep(50) + + # Stop triggers a flush via terminate/2. + GenServer.stop(pid1) + + state_file = Path.join(tmp, "watcher.state.json") + assert File.exists?(state_file), "terminate must persist watcher state" + + # New process: should rehydrate. + {:ok, _pid2} = Watcher.start_link([]) + + counts_m5 = Watcher.counts(:m5) + assert Map.get(counts_m5, [:hypatia, :scan, :complete]) == 1 + end + end + + describe "subscribe/1 (SSE fan-out)" do + setup do + case Process.whereis(Hypatia.Watcher.PubSub) do + nil -> + {:ok, pid} = Registry.start_link(keys: :duplicate, name: Hypatia.Watcher.PubSub) + on_exit(fn -> if Process.alive?(pid), do: GenServer.stop(pid) end) + + _ -> + :ok + end + + :ok + end + + test "all-events subscriber receives every kind" do + Watcher.subscribe() + + T.scan_complete(11, 2, path: "/tmp/x", severity_floor: "low") + assert_receive {:hypatia_event, [:hypatia, :scan, :complete], _measurements, _md, _ts}, 200 + + T.outcome_recorded(recipe_id: "r", repo: "x", outcome: "success", verification: "verified") + assert_receive {:hypatia_event, [:hypatia, :outcome, :recorded], _, _, _}, 200 + end + + test "filtered subscriber only sees its events" do + Watcher.subscribe(events: [[:hypatia, :verification, :result]]) + + T.scan_complete(11, 2, path: "/tmp/x", severity_floor: "low") + refute_receive {:hypatia_event, [:hypatia, :scan, :complete], _, _, _}, 100 + + T.verification_result(recipe_id: "r", repo: "x", verdict: :verified) + assert_receive {:hypatia_event, [:hypatia, :verification, :result], _, _, _}, 200 + end + + test "dead subscriber gets auto-cleaned by Registry" do + task = + Task.async(fn -> + Watcher.subscribe() + # Exit immediately after subscribing; Registry should drop us. + :done + end) + + Task.await(task) + + # Emit something; the now-dead PID's mailbox going nowhere + # must not crash the watcher. + T.scan_complete(1, 1, path: "/tmp/y", severity_floor: "low") + Process.sleep(50) + + assert Process.alive?(Process.whereis(Watcher)) + end + end end