From fd0c87c787d0626b3546fa571541c9c809711821 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 4 Dec 2025 15:17:04 +0000 Subject: add prometheus metrics --- .env.example | 20 + Cargo.lock | 45 ++- Cargo.toml | 5 + README.md | 22 +- docs/explanation/monitoring-strategy.md | 462 ---------------------- docs/explanation/monitoring.md | 99 +++++ docs/grafana/ngit-grasp-dashboard.json | 675 ++++++++++++++++++++++++++++++++ docs/how-to/prometheus-setup.md | 178 +++++++++ src/config.rs | 36 ++ src/http/mod.rs | 84 +++- src/http/nip11.rs | 6 + src/lib.rs | 1 + src/main.rs | 17 +- src/metrics/bandwidth.rs | 301 ++++++++++++++ src/metrics/connection.rs | 337 ++++++++++++++++ src/metrics/mod.rs | 469 ++++++++++++++++++++++ 16 files changed, 2285 insertions(+), 472 deletions(-) delete mode 100644 docs/explanation/monitoring-strategy.md create mode 100644 docs/explanation/monitoring.md create mode 100644 docs/grafana/ngit-grasp-dashboard.json create mode 100644 docs/how-to/prometheus-setup.md create mode 100644 src/metrics/bandwidth.rs create mode 100644 src/metrics/connection.rs create mode 100644 src/metrics/mod.rs diff --git a/.env.example b/.env.example index 0a93b1f..9207dc4 100644 --- a/.env.example +++ b/.env.example @@ -71,6 +71,26 @@ # are automatically set to temporary directories for ephemeral testing. # NGIT_DATABASE_BACKEND=lmdb +# ============================================================================ +# METRICS +# ============================================================================ + +# Enable Prometheus metrics endpoint at /metrics +# CLI: --metrics-enabled +# Default: true +# NGIT_METRICS_ENABLED=true + +# Connections per IP before flagging as potential abuse in metrics +# (display only, no rate limiting - purely for monitoring visibility) +# CLI: --metrics-connection-per-ip-abuse-threshold +# Default: 10 +# NGIT_METRICS_CONNECTION_PER_IP_ABUSE_THRESHOLD=10 + +# Number of top bandwidth repositories to track in metrics +# CLI: --metrics-top-n-repos +# Default: 10 +# NGIT_METRICS_TOP_N_REPOS=10 + # ============================================================================ # LOGGING # ============================================================================ diff --git a/Cargo.lock b/Cargo.lock index 2935855..a035f75 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -467,6 +467,19 @@ dependencies = [ "typenum", ] +[[package]] +name = "dashmap" +version = "5.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "978747c1d849a7d2ee5e8adc0159961c48fb7e5db2f06af6723b80123bb53856" +dependencies = [ + "cfg-if", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "data-encoding" version = "2.9.0" @@ -811,6 +824,12 @@ dependencies = [ "tracing", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" + [[package]] name = "hashbrown" version = "0.16.0" @@ -1163,7 +1182,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6717a8d2a5a929a1a2eb43a12812498ed141a0bcfb7e8f7844fbdbe4303bba9f" dependencies = [ "equivalent", - "hashbrown", + "hashbrown 0.16.0", ] [[package]] @@ -1353,6 +1372,7 @@ dependencies = [ "anyhow", "base64 0.22.1", "clap", + "dashmap", "dotenvy", "flate2", "futures-util", @@ -1360,9 +1380,11 @@ dependencies = [ "http-body-util", "hyper 1.8.1", "hyper-util", + "lazy_static", "nostr-lmdb", "nostr-relay-builder", "nostr-sdk 0.44.1", + "prometheus", "serde", "serde_json", "tempfile", @@ -1786,6 +1808,27 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "prometheus" +version = "0.13.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3d33c28a30771f7f96db69893f78b857f7450d7e0237e9c8fc6427a81bae7ed1" +dependencies = [ + "cfg-if", + "fnv", + "lazy_static", + "memchr", + "parking_lot", + "protobuf", + "thiserror 1.0.69", +] + +[[package]] +name = "protobuf" +version = "2.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "106dd99e98437432fed6519dedecfade6a06a73bb7b2a1e019fdd2bee5778d94" + [[package]] name = "quote" version = "1.0.41" diff --git a/Cargo.toml b/Cargo.toml index a42125d..80fa317 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,6 +28,11 @@ futures-util = "0.3" base64 = "0.22" flate2 = "1.0" +# Metrics +prometheus = "0.13" +dashmap = "5" +lazy_static = "1.4" + # Serialization serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" diff --git a/README.md b/README.md index 980526c..b63c5b1 100644 --- a/README.md +++ b/README.md @@ -94,9 +94,27 @@ Why this is useful: dedupe git data = shared object database or (GIT_ALTERNATE_OBJECT_DIRECTORIES or .git/objects/info/alternates) -### Logging +### Monitoring -See [docs/explanation/monitoring-strategy.md](docs/monitoring-strategy.md). Note: this needs refinement. +ngit-grasp exposes Prometheus metrics at `/metrics` for connection tracking, Git operations, and Nostr events. + +**Configuration Options:** + +| Option | CLI Flag | Environment Variable | Default | +|--------|----------|---------------------|---------| +| Metrics enabled | `--metrics-enabled` | `NGIT_METRICS_ENABLED` | `true` | +| Connection abuse threshold | `--metrics-connection-per-ip-abuse-threshold` | `NGIT_METRICS_CONNECTION_PER_IP_ABUSE_THRESHOLD` | `10` | +| Top N repos | `--metrics-top-n-repos` | `NGIT_METRICS_TOP_N_REPOS` | `10` | + +**Key Metrics:** +- WebSocket connections (active, unique IPs, flagged abusers) +- Git operations (clone/fetch/push rates, bandwidth, authorization results) +- Nostr events (received, stored, rejected by kind) +- Top N repositories by bandwidth + +**Privacy:** IP addresses are never exposed in metrics - only aggregate counts. + +See [Monitoring Overview](docs/explanation/monitoring.md) and [Prometheus Setup Guide](docs/how-to/prometheus-setup.md) for deployment. ### Delete Events diff --git a/docs/explanation/monitoring-strategy.md b/docs/explanation/monitoring-strategy.md deleted file mode 100644 index 4668305..0000000 --- a/docs/explanation/monitoring-strategy.md +++ /dev/null @@ -1,462 +0,0 @@ -# Monitoring Strategy - Design Document - -## Overview - -This document describes the logging and monitoring strategy for ngit-grasp, designed to help administrators: - -1. Monitor WebSocket connections per unique IP -2. Correlate resource spikes (memory, CPU) with usage patterns -3. Detect potential abuse (too many connections from single IP) -4. Support future load-based scheduling of background jobs (GRASP-02 sync) - -## Architecture - -```mermaid -flowchart TB - subgraph ngit-grasp - HTTP[HTTP Service] - WS[WebSocket Handler] - GIT[Git Handlers] - RELAY[Nostr Relay] - - subgraph Metrics Module - REG[Prometheus Registry] - CT[ConnectionTracker] - MC[Metric Counters] - end - - ME[/metrics endpoint] - end - - subgraph External - PROM[Prometheus Server] - GRAF[Grafana] - ADMIN[Admin Browser] - end - - HTTP --> ME - WS --> CT - WS --> MC - GIT --> MC - RELAY --> MC - - CT --> REG - MC --> REG - REG --> ME - - PROM -->|scrape /metrics| ME - GRAF -->|query| PROM - ADMIN -->|view dashboards| GRAF -``` - -## Metric Categories - -### 1. WebSocket Connection Metrics - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_websocket_connections_total` | Counter | - | Total WebSocket connections since startup | -| `ngit_websocket_connections_active` | Gauge | - | Current active WebSocket connections | -| `ngit_websocket_unique_ips` | Gauge | - | Number of unique IP addresses connected (NOT the IPs themselves) | -| `ngit_websocket_flagged_abusers` | Gauge | - | Number of IPs exceeding connection threshold | -| `ngit_websocket_connection_duration_seconds` | Histogram | - | Duration of WebSocket connections | -| `ngit_websocket_messages_received_total` | Counter | `type` | Messages received (REQ, EVENT, CLOSE) | -| `ngit_websocket_messages_sent_total` | Counter | `type` | Messages sent (EVENT, EOSE, OK, NOTICE) | - -**Privacy Note:** IP addresses are NEVER exposed in metrics. The `ConnectionTracker` maintains per-IP counts internally only for abuse detection, logging warnings when thresholds are exceeded. - -### 2. Git Operation Metrics - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_git_operations_total` | Counter | `operation`, `status` | Git operations (clone, fetch, push) | -| `ngit_git_operation_duration_seconds` | Histogram | `operation` | Duration of git operations | -| `ngit_git_bytes_total` | Counter | `direction` | Total bytes in/out for git operations | -| `ngit_git_push_authorization_total` | Counter | `result` | Push auth results (allowed, denied, error) | - -### 3. Top-N Repository Bandwidth Tracking - -To identify high-bandwidth repositories without creating cardinality explosion (which doesn't scale to 1000+ repos), we use a hybrid approach: - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_git_top_repos_bytes` | Gauge | `repo` | Top 10 repositories by bandwidth (refreshed every 60s) | - -**How it works:** -- All per-repo bandwidth is tracked internally in a `HashMap` -- Every 60 seconds, the top 10 are calculated and exposed to Prometheus -- Previous repo labels are cleared before setting new ones -- Prometheus only ever sees ~10 label values, keeping cardinality low - -```rust -struct BandwidthTracker { - // Internal: tracks ALL repos (memory only, not exposed) - all_repos: DashMap, - - // Exposed to Prometheus: only top 10 - top_repos_gauge: GaugeVec, - - // Refresh interval - last_refresh: Instant, -} - -impl BandwidthTracker { - fn record_transfer(&self, repo_id: &str, bytes: u64) { - self.all_repos - .entry(repo_id.to_string()) - .and_modify(|v| *v += bytes) - .or_insert(bytes); - } - - fn maybe_refresh_top_n(&self) { - if self.last_refresh.elapsed() > Duration::from_secs(60) { - self.refresh_top_n(); - } - } - - fn refresh_top_n(&self) { - let mut sorted: Vec<_> = self.all_repos.iter() - .map(|r| (r.key().clone(), *r.value())) - .collect(); - sorted.sort_by(|a, b| b.1.cmp(&a.1)); - - // Clear old labels, set new top 10 - self.top_repos_gauge.reset(); - for (repo, bytes) in sorted.into_iter().take(10) { - self.top_repos_gauge - .with_label_values(&[&repo]) - .set(bytes as i64); - } - } -} -``` - -### 4. Nostr Event Metrics - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_events_received_total` | Counter | `kind` | Events received by kind | -| `ngit_events_stored_total` | Counter | `kind` | Events successfully stored | -| `ngit_events_rejected_total` | Counter | `kind`, `reason` | Events rejected and why | - -### 5. Repository Metrics - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_repositories_total` | Gauge | - | Total repositories hosted | - -### 6. System Health Metrics - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_uptime_seconds` | Counter | - | Seconds since startup | -| `ngit_build_info` | Gauge | `version`, `commit` | Build information | - -### 7. Future: Sync Metrics (GRASP-02) - -| Metric Name | Type | Labels | Description | -|------------|------|--------|-------------| -| `ngit_sync_events_received_total` | Counter | `source` | Events from sync (live vs catchup) | -| `ngit_sync_relay_connections_active` | Gauge | - | Active outbound relay connections | -| `ngit_sync_catchup_gap_total` | Counter | - | Events found during catchup (sync failures) | - -## Connection Tracker Design - -The `ConnectionTracker` maintains per-IP connection counts internally for abuse detection. **IP addresses are never exposed in metrics** - only aggregate counts. - -```mermaid -flowchart LR - subgraph ConnectionTracker - HM[Internal: HashMap IP to Count] - TH[Abuse Threshold] - CNT[Exposed: Unique IP Count] - FLAG[Exposed: Abuse Flag Count] - end - - CONN[New Connection] --> CHECK{Count >= Threshold?} - CHECK -->|No| INC[Increment Count] - CHECK -->|Yes| FLAG_IT[Flag as Abuse] - FLAG_IT --> LOG[Log Warning - IP in log only] - FLAG_IT --> FLAG - - DISC[Disconnection] --> DEC[Decrement Count] - DEC --> CLEAN{Count == 0?} - CLEAN -->|Yes| RM[Remove from Map] - - HM --> CNT -``` - -### Data Structure - -```rust -pub struct ConnectionTracker { - /// Active connections per IP (INTERNAL ONLY - never exposed to metrics) - connections: DashMap, - /// Threshold for abuse flagging - abuse_threshold: u32, - /// Prometheus gauges (aggregate counts only, no IPs) - active_connections: IntGauge, // Total connections - unique_ips: IntGauge, // len() of HashMap - flagged_abusers: IntGauge, // Count where flagged_as_abuse == true -} - -struct ConnectionInfo { - count: u32, - first_seen: Instant, - flagged_as_abuse: bool, -} -``` - -### What Gets Exposed vs Internal - -| Data | Location | Exposed? | -|------|----------|----------| -| Total connections | Prometheus | ✅ Yes | -| Unique IP count | Prometheus | ✅ Yes | -| Flagged abuser count | Prometheus | ✅ Yes | -| Actual IP addresses | Internal HashMap | ❌ No | -| IP + abuse flag | Logs (when flagged) | ⚠️ Logs only | - -### Thread Safety - -Using `DashMap` for lock-free concurrent access, as connection tracking happens across multiple tokio tasks. - -## /metrics Endpoint - -The `/metrics` endpoint returns Prometheus text format: - -``` -# HELP ngit_websocket_connections_active Current active WebSocket connections -# TYPE ngit_websocket_connections_active gauge -ngit_websocket_connections_active 23 - -# HELP ngit_websocket_connections_by_ip Active connections per IP -# TYPE ngit_websocket_connections_by_ip gauge -ngit_websocket_connections_by_ip{ip="192.168.1.100"} 2 -ngit_websocket_connections_by_ip{ip="10.0.0.50"} 5 - -# HELP ngit_git_operations_total Git operations by type and status -# TYPE ngit_git_operations_total counter -ngit_git_operations_total{operation="clone",status="success"} 1247 -ngit_git_operations_total{operation="push",status="denied"} 12 -``` - -## Integration Points - -### HTTP Service Integration - -In [`src/http/mod.rs`](../../src/http/mod.rs): - -```rust -// Add to HttpService -struct HttpService { - // ... existing fields ... - metrics: Arc, -} - -// Add /metrics route handling -if path == "/metrics" { - let metrics_output = self.metrics.render(); - return Ok(Response::builder() - .status(200) - .header("content-type", "text/plain; version=0.0.4") - .body(Full::new(Bytes::from(metrics_output))) - .unwrap()); -} -``` - -### WebSocket Connection Tracking - -In the WebSocket upgrade handler: - -```rust -// On connection -let ip = addr.ip(); -metrics.connection_tracker.on_connect(ip); - -// Spawn connection handler -tokio::spawn(async move { - // ... handle connection ... - // On disconnect - metrics.connection_tracker.on_disconnect(ip); -}); -``` - -### Git Handler Integration - -In [`src/git/handlers.rs`](../../src/git/handlers.rs): - -```rust -// Wrap git operations with metrics -let timer = metrics.git_operation_duration.start_timer(); -let result = git::handlers::handle_upload_pack(repo_path, body_bytes).await; -timer.observe_duration(); - -metrics.git_operations_total - .with_label_values(&["clone", result_status]) - .inc(); -``` - -## Configuration - -New configuration options in [`src/config.rs`](../../src/config.rs): - -| Option | CLI Flag | Environment Variable | Default | Description | -|--------|----------|---------------------|---------|-------------| -| Metrics enabled | `--metrics-enabled` | `NGIT_METRICS_ENABLED` | `true` | Enable /metrics endpoint | -| Abuse threshold | `--abuse-threshold` | `NGIT_ABUSE_THRESHOLD` | `10` | Max connections per IP before flagging | -| Metrics path | `--metrics-path` | `NGIT_METRICS_PATH` | `/metrics` | Path for metrics endpoint | - -## Crate Dependencies - -Add to `Cargo.toml`: - -```toml -# Metrics -prometheus = "0.13" -dashmap = "5" # Lock-free concurrent HashMap -lazy_static = "1.4" # For static metric registration -``` - -## Module Structure - -``` -src/ -├── metrics/ -│ ├── mod.rs # Module exports, Metrics struct -│ ├── connection.rs # ConnectionTracker implementation -│ ├── definitions.rs # Metric definitions (lazy_static!) -│ └── render.rs # Prometheus format rendering -├── http/ -│ └── mod.rs # Add /metrics route -└── ... -``` - -## Grafana Dashboard - -A pre-built Grafana dashboard will be provided at `docs/grafana/ngit-grasp-dashboard.json` with panels for: - -1. **Overview Row** - - Active connections (gauge) - - Requests per second (graph) - - Git operations per minute (graph) - -2. **Connections Row** - - Active connections over time - - Connections by IP (top 10) - - Flagged abuse IPs (table) - -3. **Git Operations Row** - - Clone/fetch/push rates - - Push authorization results (pie chart) - - Operation duration percentiles - -4. **Events Row** - - Events received by kind - - Events rejected by reason - - Active subscriptions - -## Deployment: Prometheus on NixOS - -Example NixOS configuration for Prometheus: - -```nix -services.prometheus = { - enable = true; - scrapeConfigs = [ - { - job_name = "ngit-grasp"; - static_configs = [{ - targets = [ "localhost:8080" ]; # ngit-grasp bind address - }]; - scrape_interval = "15s"; - metrics_path = "/metrics"; - } - ]; -}; - -services.grafana = { - enable = true; - settings.server.http_port = 3000; - provision.datasources.settings.datasources = [{ - name = "Prometheus"; - type = "prometheus"; - url = "http://localhost:9090"; - }]; -}; -``` - -## Future: Load-Based Sync Scheduling - -The metrics infrastructure enables future load-based scheduling for GRASP-02 sync jobs: - -```mermaid -flowchart TD - SYNC[Sync Manager] --> CHECK{Check Load} - CHECK --> MET[Query Metrics] - MET --> CPU{CPU > 80%?} - CPU -->|Yes| DELAY[Delay 5 min] - CPU -->|No| CONN{Connections > N?} - CONN -->|Yes| DELAY - CONN -->|No| RUN[Run Sync Job] - DELAY --> CHECK -``` - -The `Metrics` struct will expose a method for checking load: - -```rust -impl Metrics { - /// Check if system is under high load - pub fn is_high_load(&self) -> bool { - let active = self.websocket_connections_active.get(); - active > self.config.high_load_threshold - } -} -``` - -## Future Enhancement: Loki for Detailed Logging - -For detailed per-repository investigation at scale, consider adding **Loki** (log aggregation) in a future iteration: - -```rust -// Structured logging with tracing -tracing::info!( - repo = %repo_id, - npub = %npub, - bytes = bytes_transferred, - operation = "clone", - duration_ms = elapsed.as_millis(), - "git_transfer_complete" -); -``` - -Loki query examples: -```logql -# Find all transfers > 10MB -{job="ngit-grasp"} |= "git_transfer_complete" | json | bytes > 10000000 - -# Sum bytes by repo in last hour -sum by (repo) ( - {job="ngit-grasp"} |= "git_transfer_complete" | json | unwrap bytes -) -``` - -This pairs with Prometheus for long-term trends while enabling ad-hoc deep dives. - -## Privacy Considerations - -- IP addresses are stored only in memory (not logged to disk by default) -- Per-IP metrics can be disabled via configuration -- Consider IP anonymization for GDPR compliance if needed - -## Summary - -| Component | Purpose | -|-----------|---------| -| `Metrics` struct | Central registry and access point | -| `ConnectionTracker` | Per-IP tracking with abuse detection | -| `/metrics` endpoint | Prometheus scraping interface | -| Grafana dashboard | Visualization and analysis | -| NixOS config | Easy deployment for operators | - -This strategy provides comprehensive observability without requiring a separate database - Prometheus handles all time-series storage and Grafana provides the visualization layer. \ No newline at end of file diff --git a/docs/explanation/monitoring.md b/docs/explanation/monitoring.md new file mode 100644 index 0000000..3b1b1ac --- /dev/null +++ b/docs/explanation/monitoring.md @@ -0,0 +1,99 @@ +# Monitoring + +ngit-grasp exposes Prometheus metrics at `/metrics` for monitoring WebSocket connections, Git operations, Nostr events, and system health. + +## Architecture + +```mermaid +flowchart TB + subgraph ngit-grasp + HTTP[HTTP Service] + WS[WebSocket Handler] + GIT[Git Handlers] + RELAY[Nostr Relay] + + subgraph Metrics Module + REG[Prometheus Registry] + CT[ConnectionTracker] + MC[Metric Counters] + end + + ME[/metrics endpoint] + end + + subgraph External + PROM[Prometheus Server] + GRAF[Grafana] + ADMIN[Admin Browser] + end + + HTTP --> ME + WS --> CT + WS --> MC + GIT --> MC + RELAY --> MC + + CT --> REG + MC --> REG + REG --> ME + + PROM -->|scrape /metrics| ME + GRAF -->|query| PROM + ADMIN -->|view dashboards| GRAF +``` + +## Configuration + +| Option | CLI Flag | Environment Variable | Default | Description | +|--------|----------|---------------------|---------|-------------| +| Metrics enabled | `--metrics-enabled` | `NGIT_METRICS_ENABLED` | `true` | Enable /metrics endpoint | +| Abuse threshold | `--abuse-threshold` | `NGIT_ABUSE_THRESHOLD` | `10` | Max connections per IP before flagging | +| Top N repos | `--top-n-repos` | `NGIT_TOP_N_REPOS` | `10` | Number of top bandwidth repos to track | + +## Privacy Model + +IP addresses are **never exposed in Prometheus metrics**. The connection tracker maintains per-IP counts internally only for abuse detection: + +| Data | Exposed in Metrics? | +|------|---------------------| +| Total connections | ✅ Yes | +| Unique IP count | ✅ Yes | +| Flagged abuser count | ✅ Yes | +| Actual IP addresses | ❌ No (internal only) | +| IP + abuse flag | ⚠️ Logs only (when flagged) | + +When an IP exceeds the abuse threshold, a warning is logged but the IP is never exposed via Prometheus. + +## Deployment + +See [Prometheus Setup Guide](../how-to/prometheus-setup.md) for NixOS configuration and Grafana dashboard provisioning. + +## Future: Load-Based Sync Scheduling (GRASP-02) + +The metrics infrastructure enables future load-based scheduling for GRASP-02 sync jobs: + +```mermaid +flowchart TD + SYNC[Sync Manager] --> CHECK{Check Load} + CHECK --> MET[Query Metrics] + MET --> CONN{Connections > N?} + CONN -->|Yes| DELAY[Delay 5 min] + CONN -->|No| RUN[Run Sync Job] + DELAY --> CHECK +``` + +## Future: Loki for Detailed Logging + +For detailed per-repository investigation at scale, consider adding **Loki** (log aggregation): + +- Structured logging with tracing crate already in place +- Loki queries enable ad-hoc deep dives (e.g., find all transfers > 10MB) +- Pairs with Prometheus for long-term trends + +## Future: Sync Metrics (GRASP-02) + +When GRASP-02 proactive sync is implemented, additional metrics will track: + +- Events received from sync (live vs catchup) +- Active outbound relay connections +- Catchup gap (events found during catchup indicating sync failures) \ No newline at end of file diff --git a/docs/grafana/ngit-grasp-dashboard.json b/docs/grafana/ngit-grasp-dashboard.json new file mode 100644 index 0000000..bd1b6fe --- /dev/null +++ b/docs/grafana/ngit-grasp-dashboard.json @@ -0,0 +1,675 @@ +{ + "annotations": { + "list": [] + }, + "editable": true, + "fiscalYearStartMonth": 0, + "graphTooltip": 0, + "id": null, + "links": [], + "liveNow": false, + "panels": [ + { + "collapsed": false, + "gridPos": { "h": 1, "w": 24, "x": 0, "y": 0 }, + "id": 1, + "title": "Overview", + "type": "row" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { "color": "green", "value": null }, + { "color": "yellow", "value": 50 }, + { "color": "red", "value": 100 } + ] + }, + "unit": "short" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 0, "y": 1 }, + "id": 2, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "auto" + }, + "pluginVersion": "10.0.0", + "targets": [ + { + "expr": "ngit_websocket_connections_active", + "legendFormat": "Active", + "refId": "A" + } + ], + "title": "Active Connections", + "type": "stat" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { "color": "green", "value": null } + ] + }, + "unit": "short" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 4, "y": 1 }, + "id": 3, + "options": { + "colorMode": "value", + "graphMode": "none", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "auto" + }, + "targets": [ + { + "expr": "ngit_websocket_unique_ips", + "legendFormat": "Unique IPs", + "refId": "A" + } + ], + "title": "Unique IPs", + "type": "stat" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { "color": "green", "value": null }, + { "color": "red", "value": 1 } + ] + }, + "unit": "short" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 8, "y": 1 }, + "id": 4, + "options": { + "colorMode": "value", + "graphMode": "none", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "auto" + }, + "targets": [ + { + "expr": "ngit_websocket_flagged_abusers", + "legendFormat": "Flagged", + "refId": "A" + } + ], + "title": "Flagged Abusers", + "type": "stat" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [{ "color": "blue", "value": null }] + }, + "unit": "short" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 12, "y": 1 }, + "id": 5, + "options": { + "colorMode": "value", + "graphMode": "none", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "auto" + }, + "targets": [ + { + "expr": "ngit_repositories_total", + "legendFormat": "Repos", + "refId": "A" + } + ], + "title": "Total Repositories", + "type": "stat" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [{ "color": "green", "value": null }] + }, + "unit": "s" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 16, "y": 1 }, + "id": 6, + "options": { + "colorMode": "value", + "graphMode": "none", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "auto" + }, + "targets": [ + { + "expr": "ngit_uptime_seconds", + "legendFormat": "Uptime", + "refId": "A" + } + ], + "title": "Uptime", + "type": "stat" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "thresholds" }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [{ "color": "purple", "value": null }] + }, + "unit": "short" + } + }, + "gridPos": { "h": 4, "w": 4, "x": 20, "y": 1 }, + "id": 7, + "options": { + "colorMode": "value", + "graphMode": "none", + "justifyMode": "auto", + "orientation": "auto", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "textMode": "value_and_name" + }, + "targets": [ + { + "expr": "ngit_build_info", + "legendFormat": "{{version}}", + "refId": "A" + } + ], + "title": "Version", + "type": "stat" + }, + { + "collapsed": false, + "gridPos": { "h": 1, "w": 24, "x": 0, "y": 5 }, + "id": 10, + "title": "WebSocket Connections", + "type": "row" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "line", + "fillOpacity": 10, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "none" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "short" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 0, "y": 6 }, + "id": 11, + "options": { + "legend": { "calcs": [], "displayMode": "list", "placement": "bottom", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "ngit_websocket_connections_active", + "legendFormat": "Active Connections", + "refId": "A" + }, + { + "expr": "ngit_websocket_unique_ips", + "legendFormat": "Unique IPs", + "refId": "B" + } + ], + "title": "Connections Over Time", + "type": "timeseries" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "line", + "fillOpacity": 10, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "none" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "short" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 12, "y": 6 }, + "id": 12, + "options": { + "legend": { "calcs": ["sum"], "displayMode": "table", "placement": "right", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "rate(ngit_websocket_messages_received_total[5m])", + "legendFormat": "Received: {{type}}", + "refId": "A" + }, + { + "expr": "rate(ngit_websocket_messages_sent_total[5m])", + "legendFormat": "Sent: {{type}}", + "refId": "B" + } + ], + "title": "Message Rate (5m)", + "type": "timeseries" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { "hideFrom": { "legend": false, "tooltip": false, "viz": false } }, + "mappings": [], + "unit": "s" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 0, "y": 14 }, + "id": 13, + "options": { + "legend": { "displayMode": "list", "placement": "right", "showLegend": true }, + "pieType": "pie", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "tooltip": { "mode": "single", "sort": "none" } + }, + "targets": [ + { + "expr": "histogram_quantile(0.5, rate(ngit_websocket_connection_duration_seconds_bucket[1h]))", + "legendFormat": "p50", + "refId": "A" + }, + { + "expr": "histogram_quantile(0.95, rate(ngit_websocket_connection_duration_seconds_bucket[1h]))", + "legendFormat": "p95", + "refId": "B" + }, + { + "expr": "histogram_quantile(0.99, rate(ngit_websocket_connection_duration_seconds_bucket[1h]))", + "legendFormat": "p99", + "refId": "C" + } + ], + "title": "Connection Duration Percentiles", + "type": "piechart" + }, + { + "collapsed": false, + "gridPos": { "h": 1, "w": 24, "x": 0, "y": 22 }, + "id": 20, + "title": "Git Operations", + "type": "row" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "bars", + "fillOpacity": 50, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "normal" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "ops" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 0, "y": 23 }, + "id": 21, + "options": { + "legend": { "calcs": ["sum"], "displayMode": "table", "placement": "right", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "rate(ngit_git_operations_total{status=\"success\"}[5m])", + "legendFormat": "{{operation}} (success)", + "refId": "A" + }, + { + "expr": "rate(ngit_git_operations_total{status=\"error\"}[5m])", + "legendFormat": "{{operation}} (error)", + "refId": "B" + } + ], + "title": "Git Operations Rate (5m)", + "type": "timeseries" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "line", + "fillOpacity": 10, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "none" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "bytes" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 12, "y": 23 }, + "id": 22, + "options": { + "legend": { "calcs": ["sum"], "displayMode": "table", "placement": "right", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "rate(ngit_git_bytes_total[5m])", + "legendFormat": "{{direction}}", + "refId": "A" + } + ], + "title": "Git Bandwidth (5m)", + "type": "timeseries" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { "hideFrom": { "legend": false, "tooltip": false, "viz": false } }, + "mappings": [], + "unit": "short" + }, + "overrides": [ + { + "matcher": { "id": "byName", "options": "denied" }, + "properties": [{ "id": "color", "value": { "fixedColor": "red", "mode": "fixed" } }] + }, + { + "matcher": { "id": "byName", "options": "allowed" }, + "properties": [{ "id": "color", "value": { "fixedColor": "green", "mode": "fixed" } }] + } + ] + }, + "gridPos": { "h": 8, "w": 6, "x": 0, "y": 31 }, + "id": 23, + "options": { + "legend": { "displayMode": "list", "placement": "right", "showLegend": true }, + "pieType": "pie", + "reduceOptions": { "calcs": ["sum"], "fields": "", "values": false }, + "tooltip": { "mode": "single", "sort": "none" } + }, + "targets": [ + { + "expr": "increase(ngit_git_push_authorization_total[24h])", + "legendFormat": "{{result}}", + "refId": "A" + } + ], + "title": "Push Authorization (24h)", + "type": "piechart" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [{ "color": "green", "value": null }] }, + "unit": "bytes" + } + }, + "gridPos": { "h": 8, "w": 18, "x": 6, "y": 31 }, + "id": 24, + "options": { + "displayMode": "gradient", + "minVizHeight": 10, + "minVizWidth": 0, + "orientation": "horizontal", + "reduceOptions": { "calcs": ["lastNotNull"], "fields": "", "values": false }, + "showUnfilled": true, + "valueMode": "color" + }, + "targets": [ + { + "expr": "topk(10, ngit_git_top_repos_bytes)", + "legendFormat": "{{repo}}", + "refId": "A" + } + ], + "title": "Top Repositories by Bandwidth", + "type": "bargauge" + }, + { + "collapsed": false, + "gridPos": { "h": 1, "w": 24, "x": 0, "y": 39 }, + "id": 30, + "title": "Nostr Events", + "type": "row" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "line", + "fillOpacity": 10, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "none" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "short" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 0, "y": 40 }, + "id": 31, + "options": { + "legend": { "calcs": ["sum"], "displayMode": "table", "placement": "right", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "rate(ngit_events_received_total[5m])", + "legendFormat": "Kind {{kind}}", + "refId": "A" + } + ], + "title": "Events Received by Kind (5m)", + "type": "timeseries" + }, + { + "datasource": { "type": "prometheus", "uid": "${datasource}" }, + "fieldConfig": { + "defaults": { + "color": { "mode": "palette-classic" }, + "custom": { + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "drawStyle": "line", + "fillOpacity": 10, + "gradientMode": "none", + "hideFrom": { "legend": false, "tooltip": false, "viz": false }, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { "type": "linear" }, + "showPoints": "never", + "spanNulls": false, + "stacking": { "group": "A", "mode": "none" }, + "thresholdsStyle": { "mode": "off" } + }, + "mappings": [], + "thresholds": { "mode": "absolute", "steps": [] }, + "unit": "short" + } + }, + "gridPos": { "h": 8, "w": 12, "x": 12, "y": 40 }, + "id": 32, + "options": { + "legend": { "calcs": ["sum"], "displayMode": "table", "placement": "right", "showLegend": true }, + "tooltip": { "mode": "multi", "sort": "none" } + }, + "targets": [ + { + "expr": "rate(ngit_events_stored_total[5m])", + "legendFormat": "Stored: Kind {{kind}}", + "refId": "A" + }, + { + "expr": "rate(ngit_events_rejected_total[5m])", + "legendFormat": "Rejected: {{reason}}", + "refId": "B" + } + ], + "title": "Events Stored vs Rejected (5m)", + "type": "timeseries" + } + ], + "refresh": "30s", + "schemaVersion": 38, + "style": "dark", + "tags": ["ngit-grasp", "nostr", "git"], + "templating": { + "list": [ + { + "current": { "selected": false, "text": "Prometheus", "value": "Prometheus" }, + "hide": 0, + "includeAll": false, + "label": "Datasource", + "multi": false, + "name": "datasource", + "options": [], + "query": "prometheus", + "refresh": 1, + "regex": "", + "skipUrlSync": false, + "type": "datasource" + } + ] + }, + "time": { "from": "now-6h", "to": "now" }, + "timepicker": {}, + "timezone": "browser", + "title": "ngit-grasp", + "uid": "ngit-grasp", + "version": 1, + "weekStart": "" +} \ No newline at end of file diff --git a/docs/how-to/prometheus-setup.md b/docs/how-to/prometheus-setup.md new file mode 100644 index 0000000..741255b --- /dev/null +++ b/docs/how-to/prometheus-setup.md @@ -0,0 +1,178 @@ +# Prometheus and Grafana Setup + +This guide shows how to configure Prometheus and Grafana to monitor ngit-grasp. + +## Prerequisites + +- ngit-grasp running with metrics enabled (default: `--metrics-enabled true`) +- Prometheus server +- Grafana (optional, for dashboards) + +## Verify Metrics Endpoint + +First, verify that ngit-grasp is exposing metrics: + +```bash +curl http://localhost:8080/metrics +``` + +You should see Prometheus-formatted metrics like: + +``` +# HELP ngit_websocket_connections_active Current active WebSocket connections +# TYPE ngit_websocket_connections_active gauge +ngit_websocket_connections_active 5 + +# HELP ngit_git_operations_total Git operations by type and status +# TYPE ngit_git_operations_total counter +ngit_git_operations_total{operation="clone",status="success"} 42 +``` + +## NixOS Configuration + +### Prometheus + +Add ngit-grasp as a scrape target: + +```nix +services.prometheus = { + enable = true; + scrapeConfigs = [ + { + job_name = "ngit-grasp"; + static_configs = [{ + targets = [ "localhost:8080" ]; # ngit-grasp bind address + }]; + scrape_interval = "15s"; + metrics_path = "/metrics"; + } + ]; +}; +``` + +### Grafana with Prometheus Datasource + +```nix +services.grafana = { + enable = true; + settings.server.http_port = 3000; + + provision.datasources.settings.datasources = [{ + name = "Prometheus"; + type = "prometheus"; + url = "http://localhost:9090"; + isDefault = true; + }]; + + # Optional: provision the ngit-grasp dashboard + provision.dashboards.settings.providers = [{ + name = "ngit-grasp"; + options.path = "/path/to/ngit-grasp/docs/grafana"; + }]; +}; +``` + +## Docker Compose Configuration + +For non-NixOS deployments: + +```yaml +version: '3.8' +services: + prometheus: + image: prom/prometheus:latest + volumes: + - ./prometheus.yml:/etc/prometheus/prometheus.yml + ports: + - "9090:9090" + + grafana: + image: grafana/grafana:latest + ports: + - "3000:3000" + volumes: + - ./docs/grafana:/var/lib/grafana/dashboards + environment: + - GF_DASHBOARDS_DEFAULT_HOME_DASHBOARD_PATH=/var/lib/grafana/dashboards/ngit-grasp-dashboard.json +``` + +With `prometheus.yml`: + +```yaml +global: + scrape_interval: 15s + +scrape_configs: + - job_name: 'ngit-grasp' + static_configs: + - targets: ['host.docker.internal:8080'] # or your ngit-grasp host + metrics_path: /metrics +``` + +## Import Dashboard + +1. Open Grafana at `http://localhost:3000` +2. Go to **Dashboards** → **Import** +3. Upload `docs/grafana/ngit-grasp-dashboard.json` +4. Select your Prometheus datasource +5. Click **Import** + +## Key Metrics to Monitor + +### Connection Health +- `ngit_websocket_connections_active` - Current active connections +- `ngit_websocket_unique_ips` - Number of unique client IPs +- `ngit_websocket_flagged_abusers` - IPs exceeding connection threshold + +### Git Operations +- `ngit_git_operations_total` - Operations by type (clone/fetch/push) and status +- `ngit_git_bytes_total` - Bandwidth by direction (in/out) +- `ngit_git_top_repos_bytes` - Top N repositories by bandwidth + +### Nostr Events +- `ngit_events_received_total` - Events received by kind +- `ngit_events_stored_total` - Events successfully stored +- `ngit_events_rejected_total` - Events rejected by reason + +### System +- `ngit_uptime_seconds` - Server uptime +- `ngit_build_info` - Version and commit info +- `ngit_repositories_total` - Total hosted repositories + +## Example Alerts + +Add to your Prometheus alerting rules: + +```yaml +groups: + - name: ngit-grasp + rules: + - alert: HighConnectionCount + expr: ngit_websocket_connections_active > 100 + for: 5m + labels: + severity: warning + annotations: + summary: "High number of WebSocket connections" + + - alert: AbusiveIPs + expr: ngit_websocket_flagged_abusers > 0 + for: 1m + labels: + severity: warning + annotations: + summary: "{{ $value }} IPs flagged for excessive connections" + + - alert: PushAuthorizationFailures + expr: rate(ngit_git_operations_total{operation="push",status="denied"}[5m]) > 0.1 + for: 5m + labels: + severity: info + annotations: + summary: "Elevated push authorization failures" +``` + +## See Also + +- [Monitoring Overview](../explanation/monitoring.md) - Architecture and design +- [Configuration Reference](../reference/configuration.md) - All config options \ No newline at end of file diff --git a/src/config.rs b/src/config.rs index d095178..025e020 100644 --- a/src/config.rs +++ b/src/config.rs @@ -71,6 +71,18 @@ pub struct Config { /// Database backend type #[arg(long, env = "NGIT_DATABASE_BACKEND", value_enum, default_value_t = DatabaseBackend::Lmdb)] pub database_backend: DatabaseBackend, + + /// Enable Prometheus metrics endpoint + #[arg(long, env = "NGIT_METRICS_ENABLED", default_value_t = true)] + pub metrics_enabled: bool, + + /// Connections per IP before flagging as potential abuse in metrics (display only, no rate limiting) + #[arg(long = "metrics-connection-per-ip-abuse-threshold", env = "NGIT_METRICS_CONNECTION_PER_IP_ABUSE_THRESHOLD", default_value_t = 10)] + pub metrics_connection_per_ip_abuse_threshold: u32, + + /// Number of top bandwidth repos to track in metrics + #[arg(long = "metrics-top-n-repos", env = "NGIT_METRICS_TOP_N_REPOS", default_value_t = 10)] + pub metrics_top_n_repos: usize, } impl Config { @@ -123,6 +135,9 @@ impl Config { relay_data_path: "./test_data/relay".to_string(), bind_address: "127.0.0.1:8080".to_string(), database_backend: DatabaseBackend::Memory, + metrics_enabled: true, + metrics_connection_per_ip_abuse_threshold: 10, + metrics_top_n_repos: 10, } } } @@ -202,4 +217,25 @@ mod tests { }; assert!(config.owner_npub.is_none()); } + + #[test] + fn test_metrics_config_defaults() { + let config = Config::for_testing(); + assert!(config.metrics_enabled); + assert_eq!(config.metrics_connection_per_ip_abuse_threshold, 10); + assert_eq!(config.metrics_top_n_repos, 10); + } + + #[test] + fn test_metrics_config_custom_values() { + let config = Config { + metrics_enabled: false, + metrics_connection_per_ip_abuse_threshold: 50, + metrics_top_n_repos: 25, + ..Config::for_testing() + }; + assert!(!config.metrics_enabled); + assert_eq!(config.metrics_connection_per_ip_abuse_threshold, 50); + assert_eq!(config.metrics_top_n_repos, 25); + } } diff --git a/src/http/mod.rs b/src/http/mod.rs index 8b1f687..f584e03 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -7,6 +7,7 @@ pub mod nip11; use std::future::Future; use std::net::SocketAddr; use std::pin::Pin; +use std::sync::Arc; use base64::Engine; use http_body_util::{BodyExt, Full}; @@ -24,6 +25,7 @@ use tokio::net::TcpListener; use crate::config::Config; use crate::git; +use crate::metrics::Metrics; use crate::nostr::builder::SharedDatabase; /// CORS headers required by GRASP-01 specification (lines 40-47) @@ -90,6 +92,8 @@ struct HttpService { remote: SocketAddr, /// Database reference for direct queries (e.g., push authorization) database: SharedDatabase, + /// Optional metrics for Prometheus endpoint + metrics: Option>, } impl HttpService { @@ -98,12 +102,14 @@ impl HttpService { config: Config, remote: SocketAddr, database: SharedDatabase, + metrics: Option>, ) -> Self { Self { relay, config, remote, database, + metrics, } } } @@ -150,6 +156,7 @@ impl Service> for HttpService { ); let repo_path = git::resolve_repo_path(&git_data_path, &npub, &identifier); + let metrics_clone = self.metrics.clone(); return Box::pin(async move { // Collect request body once before the match statement @@ -170,14 +177,31 @@ impl Service> for HttpService { .and_then(git::protocol::GitService::from_query_param); match service { - Some(svc) => git::handlers::handle_info_refs(repo_path, svc).await, + Some(svc) => { + let result = git::handlers::handle_info_refs(repo_path, svc).await; + // Track operation + if let Some(ref m) = metrics_clone { + let status = if result.is_ok() { "success" } else { "error" }; + let operation = match svc { + git::protocol::GitService::UploadPack => "fetch", + git::protocol::GitService::ReceivePack => "push", + }; + m.record_git_operation(operation, status); + } + result + } None => Err(git::handlers::GitError::RepositoryNotFound), } } // POST /git-upload-pack (clone/fetch) (m, "git-upload-pack") if m == Method::POST => { - git::handlers::handle_upload_pack(repo_path, body_bytes).await + let result = git::handlers::handle_upload_pack(repo_path, body_bytes).await; + if let Some(ref m) = metrics_clone { + let status = if result.is_ok() { "success" } else { "error" }; + m.record_git_operation("clone", status); + } + result } // POST /git-receive-pack (push) - with GRASP authorization via database @@ -187,6 +211,10 @@ impl Service> for HttpService { Ok(pk) => pk.to_hex(), Err(e) => { tracing::warn!("Invalid npub in URL {}: {}", npub, e); + // Track failed push due to invalid npub + if let Some(ref m) = metrics_clone { + m.record_git_operation("push", "error"); + } return Ok(add_cors_headers(Response::builder()) .status(hyper::StatusCode::BAD_REQUEST) .body(Full::new(Bytes::from(format!("Invalid npub: {}", e)))) @@ -194,14 +222,21 @@ impl Service> for HttpService { } }; - git::handlers::handle_receive_pack( + let result = git::handlers::handle_receive_pack( repo_path, body_bytes.clone(), Some(database.clone()), &identifier, &owner_pubkey_hex, ) - .await + .await; + + if let Some(ref m) = metrics_clone { + let status = if result.is_ok() { "success" } else { "error" }; + m.record_git_operation("push", status); + } + + result } _ => Err(git::handlers::GitError::RepositoryNotFound), @@ -333,17 +368,27 @@ impl Service> for HttpService { let addr = self.remote; let relay = self.relay.clone(); + let metrics_clone = self.metrics.clone(); tokio::spawn(async move { match hyper::upgrade::on(req).await { Ok(upgraded) => { tracing::info!("WebSocket connection established from {}", addr); + // Track connection + if let Some(ref m) = metrics_clone { + m.connection_tracker().on_connect(addr.ip()); + m.record_websocket_connection(); + } if let Err(e) = relay.take_connection(TokioIo::new(upgraded), addr).await { tracing::error!("Relay error for {}: {}", addr, e); } tracing::info!("WebSocket connection closed for {}", addr); + // Untrack connection + if let Some(ref m) = metrics_clone { + m.connection_tracker().on_disconnect(addr.ip()); + } } Err(e) => tracing::error!("Upgrade error: {}", e), } @@ -361,6 +406,33 @@ impl Service> for HttpService { } } + // Serve Prometheus metrics if enabled + if path == "/metrics" { + if let Some(ref metrics) = self.metrics { + let metrics = metrics.clone(); + return Box::pin(async move { + let output = metrics.render(); + Ok( + add_cors_headers(Response::builder().header("server", "ngit-grasp")) + .status(200) + .header("content-type", "text/plain; version=0.0.4; charset=utf-8") + .body(Full::new(Bytes::from(output))) + .unwrap(), + ) + }); + } else { + // Metrics disabled + return Box::pin(async move { + Ok( + add_cors_headers(Response::builder().header("server", "ngit-grasp")) + .status(404) + .body(Full::new(Bytes::from("Metrics disabled"))) + .unwrap(), + ) + }); + } + } + // Serve static icon at /icon.png if path == "/icon.png" { return Box::pin(async move { @@ -419,10 +491,12 @@ fn derive_accept_key(request_key: &[u8]) -> String { /// * `config` - Server configuration /// * `relay` - The LocalRelay for WebSocket connections /// * `database` - The database for direct queries (e.g., push authorization) +/// * `metrics` - Optional metrics for Prometheus endpoint pub async fn run_server( config: Config, relay: LocalRelay, database: SharedDatabase, + metrics: Option>, ) -> anyhow::Result<()> { let bind_addr: SocketAddr = config.bind_address.parse()?; @@ -435,7 +509,7 @@ pub async fn run_server( loop { let (socket, addr) = listener.accept().await?; let io = TokioIo::new(socket); - let service = HttpService::new(relay.clone(), config.clone(), addr, database.clone()); + let service = HttpService::new(relay.clone(), config.clone(), addr, database.clone(), metrics.clone()); tokio::spawn(async move { if let Err(e) = http1::Builder::new() diff --git a/src/http/nip11.rs b/src/http/nip11.rs index 1ed80de..1723601 100644 --- a/src/http/nip11.rs +++ b/src/http/nip11.rs @@ -102,6 +102,9 @@ mod tests { relay_data_path: "./data/relay".to_string(), bind_address: "127.0.0.1:8080".to_string(), database_backend: crate::config::DatabaseBackend::Memory, + metrics_enabled: true, + metrics_connection_per_ip_abuse_threshold: 10, + metrics_top_n_repos: 10, }; let doc = RelayInformationDocument::from_config(&config); @@ -133,6 +136,9 @@ mod tests { relay_data_path: "./data/relay".to_string(), bind_address: "127.0.0.1:8080".to_string(), database_backend: crate::config::DatabaseBackend::Memory, + metrics_enabled: true, + metrics_connection_per_ip_abuse_threshold: 10, + metrics_top_n_repos: 10, }; let doc = RelayInformationDocument::from_config(&config); diff --git a/src/lib.rs b/src/lib.rs index 9ccc212..4d5aab0 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,4 +1,5 @@ pub mod config; pub mod git; pub mod http; +pub mod metrics; pub mod nostr; diff --git a/src/main.rs b/src/main.rs index f80e920..9200cc2 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,10 +1,14 @@ +use std::sync::Arc; + use anyhow::Result; use tracing::{info, Level}; use tracing_subscriber::FmtSubscriber; use ngit_grasp::{ config::{Config, DatabaseBackend}, - http, nostr, + http, + metrics::Metrics, + nostr, }; #[tokio::main] @@ -29,6 +33,15 @@ async fn main() -> Result<()> { } info!("Database backend: {}", config.database_backend); + // Initialize metrics if enabled + let metrics = if config.metrics_enabled { + info!("Metrics enabled on /metrics endpoint"); + Some(Arc::new(Metrics::new(config.metrics_connection_per_ip_abuse_threshold))) + } else { + info!("Metrics disabled"); + None + }; + // Create Nostr relay with NIP-34 validation // Returns both the relay and database for direct queries in handlers if let Ok(relay_with_db) = nostr::builder::create_relay(&config) { @@ -39,7 +52,7 @@ async fn main() -> Result<()> { // Start HTTP server with integrated relay and database info!("Starting HTTP server on {}", config.bind_address); - http::run_server(config, relay_with_db.relay, relay_with_db.database).await?; + http::run_server(config, relay_with_db.relay, relay_with_db.database, metrics).await?; } Ok(()) diff --git a/src/metrics/bandwidth.rs b/src/metrics/bandwidth.rs new file mode 100644 index 0000000..d2c53e8 --- /dev/null +++ b/src/metrics/bandwidth.rs @@ -0,0 +1,301 @@ +//! Repository bandwidth tracking with cardinality control. +//! +//! This module tracks bandwidth per repository but only exposes the top N +//! repositories to Prometheus to prevent cardinality explosion with many repos. +//! +//! # Cardinality Control +//! +//! - All per-repo bandwidth is tracked internally in a `DashMap` +//! - Every 60 seconds, the top 10 are calculated and exposed to Prometheus +//! - Previous repo labels are cleared before setting new ones +//! - Prometheus only ever sees ~10 label values, keeping cardinality low + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Duration, Instant}; + +use dashmap::DashMap; +use prometheus::{GaugeVec, Opts, Registry}; + +/// Default number of top repositories to expose in metrics +const DEFAULT_TOP_N: usize = 10; + +/// Default refresh interval for top-N calculation (60 seconds) +const DEFAULT_REFRESH_INTERVAL: Duration = Duration::from_secs(60); + +/// Tracks bandwidth per repository with top-N exposure to Prometheus. +/// +/// # Design +/// +/// All repositories are tracked internally for accurate total bandwidth, +/// but only the top N by bytes transferred are exposed to Prometheus. +/// This prevents cardinality explosion when hosting thousands of repositories. +/// +/// # Thread Safety +/// +/// Uses `DashMap` for lock-free concurrent access and atomics for +/// the refresh timestamp. +pub struct BandwidthTracker { + /// Internal: tracks ALL repos (memory only, not exposed) + all_repos: DashMap, + + /// Exposed to Prometheus: only top N repos + top_repos_gauge: GaugeVec, + + /// Last refresh timestamp (stored as nanos since some epoch) + last_refresh_nanos: AtomicU64, + + /// Instant when the tracker was created (for relative timing) + start_instant: Instant, + + /// Number of top repos to expose + top_n: usize, + + /// Refresh interval + refresh_interval: Duration, +} + +impl BandwidthTracker { + /// Creates a new BandwidthTracker and registers metrics with Prometheus. + /// + /// Uses default settings: + /// - Top 10 repositories exposed + /// - 60 second refresh interval + pub fn new(registry: &Registry) -> Self { + Self::with_config(registry, DEFAULT_TOP_N, DEFAULT_REFRESH_INTERVAL) + } + + /// Creates a new BandwidthTracker with custom configuration. + /// + /// # Arguments + /// + /// * `registry` - Prometheus registry to register metrics with + /// * `top_n` - Number of top repositories to expose in metrics + /// * `refresh_interval` - How often to recalculate the top-N list + pub fn with_config(registry: &Registry, top_n: usize, refresh_interval: Duration) -> Self { + let top_repos_gauge = GaugeVec::new( + Opts::new( + "ngit_git_top_repos_bytes", + "Top repositories by bandwidth (refreshed periodically)", + ), + &["repo"], + ) + .unwrap(); + registry.register(Box::new(top_repos_gauge.clone())).unwrap(); + + Self { + all_repos: DashMap::new(), + top_repos_gauge, + last_refresh_nanos: AtomicU64::new(0), + start_instant: Instant::now(), + top_n, + refresh_interval, + } + } + + /// Records bytes transferred for a repository. + /// + /// # Arguments + /// + /// * `repo_id` - Repository identifier (e.g., npub or repo name) + /// * `bytes` - Number of bytes transferred + pub fn record_transfer(&self, repo_id: &str, bytes: u64) { + self.all_repos + .entry(repo_id.to_string()) + .and_modify(|v| *v = v.saturating_add(bytes)) + .or_insert(bytes); + } + + /// Conditionally refreshes the top-N list if the refresh interval has elapsed. + /// + /// This method is designed to be called frequently (e.g., on every + /// `/metrics` request) without performance impact - it only does work + /// when the refresh interval has elapsed. + pub fn maybe_refresh_top_n(&self) { + let elapsed_nanos = self.start_instant.elapsed().as_nanos() as u64; + let last_refresh = self.last_refresh_nanos.load(Ordering::Relaxed); + let interval_nanos = self.refresh_interval.as_nanos() as u64; + + // Check if enough time has passed since last refresh + if elapsed_nanos.saturating_sub(last_refresh) >= interval_nanos { + // Try to update the timestamp atomically to prevent concurrent refreshes + if self + .last_refresh_nanos + .compare_exchange(last_refresh, elapsed_nanos, Ordering::SeqCst, Ordering::Relaxed) + .is_ok() + { + self.refresh_top_n(); + } + } + } + + /// Forces a refresh of the top-N list. + /// + /// This recalculates which repositories are in the top N by bandwidth + /// and updates the Prometheus gauges accordingly. + pub fn refresh_top_n(&self) { + // Collect all repo data + let mut sorted: Vec<_> = self + .all_repos + .iter() + .map(|r| (r.key().clone(), *r.value())) + .collect(); + + // Sort by bytes descending + sorted.sort_by(|a, b| b.1.cmp(&a.1)); + + // Clear old labels and set new top N + self.top_repos_gauge.reset(); + for (repo, bytes) in sorted.into_iter().take(self.top_n) { + self.top_repos_gauge + .with_label_values(&[&repo]) + .set(bytes as f64); + } + } + + /// Returns the total bytes transferred for a specific repository. + /// + /// Returns `None` if the repository has not been seen. + pub fn get_repo_bytes(&self, repo_id: &str) -> Option { + self.all_repos.get(repo_id).map(|v| *v) + } + + /// Returns the total bytes transferred across all repositories. + pub fn total_bytes(&self) -> u64 { + self.all_repos.iter().map(|r| *r.value()).sum() + } + + /// Returns the number of repositories being tracked. + pub fn repo_count(&self) -> usize { + self.all_repos.len() + } + + /// Returns the top N repositories by bandwidth. + /// + /// This is a snapshot and may not match the Prometheus gauges if + /// a refresh hasn't occurred recently. + pub fn get_top_repos(&self) -> Vec<(String, u64)> { + let mut sorted: Vec<_> = self + .all_repos + .iter() + .map(|r| (r.key().clone(), *r.value())) + .collect(); + + sorted.sort_by(|a, b| b.1.cmp(&a.1)); + sorted.truncate(self.top_n); + sorted + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn test_registry() -> Registry { + Registry::new() + } + + #[test] + fn test_bandwidth_tracking() { + let registry = test_registry(); + let tracker = BandwidthTracker::new(®istry); + + // Record transfers + tracker.record_transfer("repo-a", 1000); + tracker.record_transfer("repo-b", 2000); + tracker.record_transfer("repo-a", 500); // Additional transfer to repo-a + + assert_eq!(tracker.get_repo_bytes("repo-a"), Some(1500)); + assert_eq!(tracker.get_repo_bytes("repo-b"), Some(2000)); + assert_eq!(tracker.get_repo_bytes("repo-c"), None); + assert_eq!(tracker.total_bytes(), 3500); + assert_eq!(tracker.repo_count(), 2); + } + + #[test] + fn test_top_n_repos() { + let registry = test_registry(); + let tracker = BandwidthTracker::with_config(®istry, 3, Duration::from_secs(60)); + + // Create 5 repos with different bandwidth + tracker.record_transfer("repo-1", 100); + tracker.record_transfer("repo-2", 500); + tracker.record_transfer("repo-3", 200); + tracker.record_transfer("repo-4", 800); + tracker.record_transfer("repo-5", 300); + + let top = tracker.get_top_repos(); + assert_eq!(top.len(), 3); + assert_eq!(top[0], ("repo-4".to_string(), 800)); + assert_eq!(top[1], ("repo-2".to_string(), 500)); + assert_eq!(top[2], ("repo-5".to_string(), 300)); + } + + #[test] + fn test_refresh_updates_gauge() { + let registry = test_registry(); + let tracker = BandwidthTracker::new(®istry); + + tracker.record_transfer("high-bandwidth-repo", 10_000_000); + tracker.record_transfer("low-bandwidth-repo", 1000); + + // Force a refresh + tracker.refresh_top_n(); + + // Verify the gauge values (we can't easily access them directly, + // but we can verify the tracker state is correct) + assert_eq!(tracker.repo_count(), 2); + assert_eq!(tracker.total_bytes(), 10_001_000); + } + + #[test] + fn test_saturating_add() { + let registry = test_registry(); + let tracker = BandwidthTracker::new(®istry); + + // Test that we don't overflow + tracker.record_transfer("huge-repo", u64::MAX - 100); + tracker.record_transfer("huge-repo", 200); + + // Should saturate to MAX, not overflow + assert_eq!(tracker.get_repo_bytes("huge-repo"), Some(u64::MAX)); + } + + #[test] + fn test_maybe_refresh_respects_interval() { + let registry = test_registry(); + // Use a very short interval for testing + let tracker = BandwidthTracker::with_config(®istry, 10, Duration::from_millis(10)); + + tracker.record_transfer("repo-a", 1000); + + // First call should trigger refresh (no previous refresh) + tracker.maybe_refresh_top_n(); + + // Add more data + tracker.record_transfer("repo-b", 2000); + + // Immediate second call should NOT trigger refresh + let count_before = tracker.repo_count(); + tracker.maybe_refresh_top_n(); + assert_eq!(tracker.repo_count(), count_before); + + // Wait for interval to pass + std::thread::sleep(Duration::from_millis(15)); + + // Now it should refresh + tracker.maybe_refresh_top_n(); + } + + #[test] + fn test_empty_tracker() { + let registry = test_registry(); + let tracker = BandwidthTracker::new(®istry); + + assert_eq!(tracker.total_bytes(), 0); + assert_eq!(tracker.repo_count(), 0); + assert!(tracker.get_top_repos().is_empty()); + + // Refresh should not panic on empty data + tracker.refresh_top_n(); + } +} \ No newline at end of file diff --git a/src/metrics/connection.rs b/src/metrics/connection.rs new file mode 100644 index 0000000..6a7f406 --- /dev/null +++ b/src/metrics/connection.rs @@ -0,0 +1,337 @@ +//! Connection tracking with privacy-preserving abuse detection. +//! +//! This module tracks WebSocket connections per IP address internally for abuse +//! detection, but NEVER exposes IP addresses in Prometheus metrics. Only aggregate +//! counts are exposed. +//! +//! # Privacy Model +//! +//! | Data | Location | Exposed? | +//! |------|----------|----------| +//! | Total connections | Prometheus | ✅ Yes | +//! | Unique IP count | Prometheus | ✅ Yes | +//! | Flagged abuser count | Prometheus | ✅ Yes | +//! | Actual IP addresses | Internal HashMap | ❌ No | +//! | IP + abuse flag | Logs (when flagged) | ⚠️ Logs only | + +use std::net::IpAddr; +use std::time::Instant; + +use dashmap::DashMap; +use prometheus::{IntGauge, Opts, Registry}; +use tracing::warn; + +/// Information about connections from a specific IP address. +struct ConnectionInfo { + /// Number of active connections from this IP + count: u32, + /// When the first connection from this IP was established + first_seen: Instant, + /// Whether this IP has been flagged as potentially abusive + flagged_as_abuse: bool, +} + +/// Tracks WebSocket connections per IP with abuse detection. +/// +/// # Thread Safety +/// +/// Uses `DashMap` for lock-free concurrent access, as connection tracking +/// happens across multiple tokio tasks. +/// +/// # Privacy +/// +/// IP addresses are stored internally only for abuse detection and are +/// NEVER exposed in Prometheus metrics. Only aggregate counts are exposed: +/// - Total active connections +/// - Number of unique IPs +/// - Number of IPs flagged as potential abusers +pub struct ConnectionTracker { + /// Active connections per IP (INTERNAL ONLY - never exposed to metrics) + connections: DashMap, + + /// Threshold for abuse flagging (connections per IP) + abuse_threshold: u32, + + /// Prometheus gauge: total active connections + active_connections: IntGauge, + + /// Prometheus gauge: number of unique IPs connected + unique_ips: IntGauge, + + /// Prometheus gauge: number of IPs flagged as potential abusers + flagged_abusers: IntGauge, +} + +impl ConnectionTracker { + /// Creates a new ConnectionTracker and registers metrics with Prometheus. + /// + /// # Arguments + /// + /// * `abuse_threshold` - Number of connections from a single IP before flagging + /// * `registry` - Prometheus registry to register metrics with + pub fn new(abuse_threshold: u32, registry: &Registry) -> Self { + let active_connections = IntGauge::with_opts( + Opts::new( + "ngit_websocket_connections_active", + "Current active WebSocket connections", + ) + ).unwrap(); + registry.register(Box::new(active_connections.clone())).unwrap(); + + let unique_ips = IntGauge::with_opts( + Opts::new( + "ngit_websocket_unique_ips", + "Number of unique IP addresses connected (NOT the IPs themselves)", + ) + ).unwrap(); + registry.register(Box::new(unique_ips.clone())).unwrap(); + + let flagged_abusers = IntGauge::with_opts( + Opts::new( + "ngit_websocket_flagged_abusers", + "Number of IPs exceeding connection threshold", + ) + ).unwrap(); + registry.register(Box::new(flagged_abusers.clone())).unwrap(); + + Self { + connections: DashMap::new(), + abuse_threshold, + active_connections, + unique_ips, + flagged_abusers, + } + } + + /// Called when a new WebSocket connection is established. + /// + /// This method: + /// 1. Increments the connection count for this IP + /// 2. Checks if the IP has exceeded the abuse threshold + /// 3. Logs a warning if abuse is detected (IP is logged here only) + /// 4. Updates Prometheus metrics (aggregate counts only) + /// + /// # Privacy + /// + /// The IP address is logged only when abuse is detected. It is NEVER + /// exposed in Prometheus metrics. + pub fn on_connect(&self, ip: IpAddr) { + let mut is_new_ip = false; + let mut newly_flagged = false; + + self.connections + .entry(ip) + .and_modify(|info| { + info.count += 1; + // Check if this connection pushes us over the threshold + if !info.flagged_as_abuse && info.count >= self.abuse_threshold { + info.flagged_as_abuse = true; + newly_flagged = true; + } + }) + .or_insert_with(|| { + is_new_ip = true; + ConnectionInfo { + count: 1, + first_seen: Instant::now(), + flagged_as_abuse: false, + } + }); + + // Update Prometheus metrics (aggregate counts only) + self.active_connections.inc(); + + if is_new_ip { + self.unique_ips.inc(); + } + + if newly_flagged { + self.flagged_abusers.inc(); + // Log the abuse detection - IP is only exposed in logs, not metrics + warn!( + ip = %ip, + threshold = self.abuse_threshold, + "Potential abuse detected: IP exceeded connection threshold" + ); + } + } + + /// Called when a WebSocket connection is closed. + /// + /// This method: + /// 1. Decrements the connection count for this IP + /// 2. Removes the IP from tracking if count reaches 0 + /// 3. Updates the abuse flag count if the IP was flagged + /// 4. Updates Prometheus metrics (aggregate counts only) + pub fn on_disconnect(&self, ip: IpAddr) { + let mut remove_entry = false; + let mut was_flagged = false; + let mut had_connection = false; + + if let Some(mut entry) = self.connections.get_mut(&ip) { + had_connection = true; + entry.count = entry.count.saturating_sub(1); + if entry.count == 0 { + remove_entry = true; + was_flagged = entry.flagged_as_abuse; + } + } + + // Remove the entry if count is 0 + if remove_entry { + self.connections.remove(&ip); + self.unique_ips.dec(); + if was_flagged { + self.flagged_abusers.dec(); + } + } + + // Update total connections only if this IP had a tracked connection + if had_connection { + self.active_connections.dec(); + } + } + + /// Returns the current number of active connections. + pub fn active_connections(&self) -> u64 { + self.active_connections.get() as u64 + } + + /// Returns the current number of unique IPs. + pub fn unique_ip_count(&self) -> u64 { + self.unique_ips.get() as u64 + } + + /// Returns the current number of flagged abusers. + pub fn flagged_abuser_count(&self) -> u64 { + self.flagged_abusers.get() as u64 + } + + /// Returns the connection count for a specific IP (for internal use only). + /// + /// # Privacy + /// + /// This is an internal method. The returned data should NEVER be exposed + /// in metrics or logs without privacy consideration. + #[cfg(test)] + pub(crate) fn connection_count(&self, ip: &IpAddr) -> Option { + self.connections.get(ip).map(|info| info.count) + } + + /// Returns whether an IP is flagged as abusive (for internal use only). + #[cfg(test)] + pub(crate) fn is_flagged(&self, ip: &IpAddr) -> bool { + self.connections + .get(ip) + .map(|info| info.flagged_as_abuse) + .unwrap_or(false) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::net::{Ipv4Addr, Ipv6Addr}; + + fn test_registry() -> Registry { + Registry::new() + } + + #[test] + fn test_connection_tracking() { + let registry = test_registry(); + let tracker = ConnectionTracker::new(5, ®istry); + let ip = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 1)); + + // Connect + tracker.on_connect(ip); + assert_eq!(tracker.active_connections(), 1); + assert_eq!(tracker.unique_ip_count(), 1); + assert_eq!(tracker.connection_count(&ip), Some(1)); + + // Connect again from same IP + tracker.on_connect(ip); + assert_eq!(tracker.active_connections(), 2); + assert_eq!(tracker.unique_ip_count(), 1); // Still 1 unique IP + assert_eq!(tracker.connection_count(&ip), Some(2)); + + // Disconnect one + tracker.on_disconnect(ip); + assert_eq!(tracker.active_connections(), 1); + assert_eq!(tracker.unique_ip_count(), 1); + assert_eq!(tracker.connection_count(&ip), Some(1)); + + // Disconnect last + tracker.on_disconnect(ip); + assert_eq!(tracker.active_connections(), 0); + assert_eq!(tracker.unique_ip_count(), 0); + assert_eq!(tracker.connection_count(&ip), None); + } + + #[test] + fn test_multiple_ips() { + let registry = test_registry(); + let tracker = ConnectionTracker::new(5, ®istry); + let ip1 = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 1)); + let ip2 = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 2)); + let ip3 = IpAddr::V6(Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, 1)); + + tracker.on_connect(ip1); + tracker.on_connect(ip2); + tracker.on_connect(ip3); + + assert_eq!(tracker.active_connections(), 3); + assert_eq!(tracker.unique_ip_count(), 3); + + tracker.on_disconnect(ip2); + assert_eq!(tracker.active_connections(), 2); + assert_eq!(tracker.unique_ip_count(), 2); + } + + #[test] + fn test_abuse_detection() { + let registry = test_registry(); + let threshold = 3; + let tracker = ConnectionTracker::new(threshold, ®istry); + let abuser_ip = IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)); + let normal_ip = IpAddr::V4(Ipv4Addr::new(10, 0, 0, 2)); + + // Normal user with 1 connection + tracker.on_connect(normal_ip); + assert!(!tracker.is_flagged(&normal_ip)); + assert_eq!(tracker.flagged_abuser_count(), 0); + + // Abuser approaching threshold + tracker.on_connect(abuser_ip); + tracker.on_connect(abuser_ip); + assert!(!tracker.is_flagged(&abuser_ip)); + assert_eq!(tracker.flagged_abuser_count(), 0); + + // Abuser hits threshold + tracker.on_connect(abuser_ip); + assert!(tracker.is_flagged(&abuser_ip)); + assert_eq!(tracker.flagged_abuser_count(), 1); + + // Normal user still not flagged + assert!(!tracker.is_flagged(&normal_ip)); + + // Abuser disconnects all - should be removed from flagged count + tracker.on_disconnect(abuser_ip); + tracker.on_disconnect(abuser_ip); + tracker.on_disconnect(abuser_ip); + assert_eq!(tracker.flagged_abuser_count(), 0); + assert_eq!(tracker.active_connections(), 1); // Only normal user remains + } + + #[test] + fn test_disconnect_without_connect() { + let registry = test_registry(); + let tracker = ConnectionTracker::new(5, ®istry); + let ip = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 1)); + + // Disconnect without connect should not panic or go negative + tracker.on_disconnect(ip); + assert_eq!(tracker.active_connections(), 0); + assert_eq!(tracker.unique_ip_count(), 0); + } +} \ No newline at end of file diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs new file mode 100644 index 0000000..4a4fe57 --- /dev/null +++ b/src/metrics/mod.rs @@ -0,0 +1,469 @@ +//! Prometheus metrics for ngit-grasp relay. +//! +//! This module provides comprehensive monitoring metrics including: +//! - WebSocket connection tracking (with privacy-preserving IP aggregation) +//! - Git operation metrics (clone, fetch, push) +//! - Repository bandwidth tracking (top-N only for cardinality control) +//! - Nostr event metrics +//! +//! # Privacy +//! IP addresses are NEVER exposed in metrics. The `ConnectionTracker` maintains +//! per-IP counts internally only for abuse detection. Only aggregate counts +//! are exposed to Prometheus. + +pub mod bandwidth; +pub mod connection; + +use std::sync::Arc; +use std::time::Instant; + +use lazy_static::lazy_static; +use prometheus::{ + Counter, CounterVec, Encoder, Gauge, GaugeVec, Histogram, HistogramOpts, HistogramVec, Opts, + Registry, TextEncoder, +}; + +use bandwidth::BandwidthTracker; +use connection::ConnectionTracker; + +lazy_static! { + /// Global Prometheus registry for ngit-grasp metrics + pub static ref REGISTRY: Registry = Registry::new(); +} + +/// Central metrics collection for ngit-grasp relay. +/// +/// Thread-safe and designed for concurrent access from multiple tokio tasks. +#[derive(Clone)] +pub struct Metrics { + inner: Arc, +} + +struct MetricsInner { + /// Connection tracking with abuse detection + pub connection_tracker: ConnectionTracker, + + /// Repository bandwidth tracking (top-N only) + pub bandwidth_tracker: BandwidthTracker, + + // === WebSocket Metrics === + /// Total WebSocket connections since startup + pub websocket_connections_total: Counter, + /// Connection duration histogram + pub websocket_connection_duration: Histogram, + /// Messages received by type (REQ, EVENT, CLOSE) + pub websocket_messages_received: CounterVec, + /// Messages sent by type (EVENT, EOSE, OK, NOTICE) + pub websocket_messages_sent: CounterVec, + + // === Git Operation Metrics === + /// Git operations by type and status + pub git_operations_total: CounterVec, + /// Git operation duration histogram + pub git_operation_duration: HistogramVec, + /// Total bytes transferred + pub git_bytes_total: CounterVec, + /// Push authorization results + pub git_push_authorization: CounterVec, + + // === Nostr Event Metrics === + /// Events received by kind + pub events_received_total: CounterVec, + /// Events successfully stored by kind + pub events_stored_total: CounterVec, + /// Events rejected by kind and reason + pub events_rejected_total: CounterVec, + + // === Repository Metrics === + /// Total repositories hosted + pub repositories_total: Gauge, + + // === System Health Metrics === + /// Server start time for uptime calculation + pub start_time: Instant, + /// Build information gauge + pub build_info: GaugeVec, +} + +impl Metrics { + /// Creates a new Metrics instance and registers all metrics with Prometheus. + /// + /// # Arguments + /// * `abuse_threshold` - Number of connections from a single IP before flagging as abuse + pub fn new(abuse_threshold: u32) -> Self { + let inner = MetricsInner::new(abuse_threshold); + Self { + inner: Arc::new(inner), + } + } + + /// Returns the connection tracker for WebSocket connection management. + pub fn connection_tracker(&self) -> &ConnectionTracker { + &self.inner.connection_tracker + } + + /// Returns the bandwidth tracker for repository bandwidth tracking. + pub fn bandwidth_tracker(&self) -> &BandwidthTracker { + &self.inner.bandwidth_tracker + } + + // === WebSocket Recording Methods === + + /// Record a new WebSocket connection + pub fn record_websocket_connection(&self) { + self.inner.websocket_connections_total.inc(); + } + + /// Start timing a WebSocket connection, returns timer that records on drop + pub fn start_connection_timer(&self) -> HistogramTimer { + HistogramTimer::new(self.inner.websocket_connection_duration.clone()) + } + + /// Record a received WebSocket message + pub fn record_message_received(&self, msg_type: &str) { + self.inner + .websocket_messages_received + .with_label_values(&[msg_type]) + .inc(); + } + + /// Record a sent WebSocket message + pub fn record_message_sent(&self, msg_type: &str) { + self.inner + .websocket_messages_sent + .with_label_values(&[msg_type]) + .inc(); + } + + // === Git Operation Recording Methods === + + /// Record a git operation completion + pub fn record_git_operation(&self, operation: &str, status: &str) { + self.inner + .git_operations_total + .with_label_values(&[operation, status]) + .inc(); + } + + /// Start timing a git operation, returns a timer + pub fn start_git_operation_timer(&self, operation: &str) -> GitOperationTimer { + GitOperationTimer::new(self.inner.git_operation_duration.clone(), operation.to_string()) + } + + /// Record bytes transferred for a git operation + pub fn record_git_bytes(&self, direction: &str, bytes: u64) { + self.inner + .git_bytes_total + .with_label_values(&[direction]) + .inc_by(bytes as f64); + } + + /// Record a push authorization result + pub fn record_push_authorization(&self, result: &str) { + self.inner + .git_push_authorization + .with_label_values(&[result]) + .inc(); + } + + // === Nostr Event Recording Methods === + + /// Record a received Nostr event + pub fn record_event_received(&self, kind: u64) { + self.inner + .events_received_total + .with_label_values(&[&kind.to_string()]) + .inc(); + } + + /// Record a stored Nostr event + pub fn record_event_stored(&self, kind: u64) { + self.inner + .events_stored_total + .with_label_values(&[&kind.to_string()]) + .inc(); + } + + /// Record a rejected Nostr event + pub fn record_event_rejected(&self, kind: u64, reason: &str) { + self.inner + .events_rejected_total + .with_label_values(&[&kind.to_string(), reason]) + .inc(); + } + + // === Repository Metrics === + + /// Set the total number of repositories + pub fn set_repositories_total(&self, count: u64) { + self.inner.repositories_total.set(count as f64); + } + + /// Increment the repository count + pub fn inc_repositories_total(&self) { + self.inner.repositories_total.inc(); + } + + // === Rendering === + + /// Render all metrics in Prometheus text format. + /// + /// This method: + /// 1. Refreshes the top-N bandwidth metrics if needed + /// 2. Updates uptime + /// 3. Gathers all metrics from the registry + /// 4. Encodes them in Prometheus text format + pub fn render(&self) -> String { + // Refresh top-N bandwidth repos if needed + self.inner.bandwidth_tracker.maybe_refresh_top_n(); + + // Gather and encode metrics + let encoder = TextEncoder::new(); + let metric_families = REGISTRY.gather(); + let mut buffer = Vec::new(); + encoder.encode(&metric_families, &mut buffer).unwrap(); + + // Add uptime as a comment (it's derived, not a registered metric) + let uptime = self.inner.start_time.elapsed().as_secs(); + let mut output = String::from_utf8(buffer).unwrap(); + output.push_str(&format!( + "\n# HELP ngit_uptime_seconds Seconds since server startup\n# TYPE ngit_uptime_seconds counter\nngit_uptime_seconds {}\n", + uptime + )); + + output + } + + /// Check if the system is under high load (for sync scheduling) + pub fn is_high_load(&self, threshold: u64) -> bool { + self.inner.connection_tracker.active_connections() > threshold + } +} + +impl MetricsInner { + fn new(abuse_threshold: u32) -> Self { + // Create connection tracker + let connection_tracker = ConnectionTracker::new(abuse_threshold, ®ISTRY); + + // Create bandwidth tracker + let bandwidth_tracker = BandwidthTracker::new(®ISTRY); + + // WebSocket metrics + let websocket_connections_total = Counter::with_opts( + Opts::new( + "ngit_websocket_connections_total", + "Total WebSocket connections since startup", + ) + ).unwrap(); + REGISTRY.register(Box::new(websocket_connections_total.clone())).unwrap(); + + let websocket_connection_duration = Histogram::with_opts( + HistogramOpts::new( + "ngit_websocket_connection_duration_seconds", + "Duration of WebSocket connections", + ) + .buckets(vec![1.0, 5.0, 15.0, 30.0, 60.0, 300.0, 900.0, 3600.0]), + ).unwrap(); + REGISTRY.register(Box::new(websocket_connection_duration.clone())).unwrap(); + + let websocket_messages_received = CounterVec::new( + Opts::new( + "ngit_websocket_messages_received_total", + "WebSocket messages received by type", + ), + &["type"], + ).unwrap(); + REGISTRY.register(Box::new(websocket_messages_received.clone())).unwrap(); + + let websocket_messages_sent = CounterVec::new( + Opts::new( + "ngit_websocket_messages_sent_total", + "WebSocket messages sent by type", + ), + &["type"], + ).unwrap(); + REGISTRY.register(Box::new(websocket_messages_sent.clone())).unwrap(); + + // Git operation metrics + let git_operations_total = CounterVec::new( + Opts::new( + "ngit_git_operations_total", + "Git operations by type and status", + ), + &["operation", "status"], + ).unwrap(); + REGISTRY.register(Box::new(git_operations_total.clone())).unwrap(); + + let git_operation_duration = HistogramVec::new( + HistogramOpts::new( + "ngit_git_operation_duration_seconds", + "Duration of git operations", + ) + .buckets(vec![0.1, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0]), + &["operation"], + ).unwrap(); + REGISTRY.register(Box::new(git_operation_duration.clone())).unwrap(); + + let git_bytes_total = CounterVec::new( + Opts::new( + "ngit_git_bytes_total", + "Total bytes transferred for git operations", + ), + &["direction"], + ).unwrap(); + REGISTRY.register(Box::new(git_bytes_total.clone())).unwrap(); + + let git_push_authorization = CounterVec::new( + Opts::new( + "ngit_git_push_authorization_total", + "Push authorization results", + ), + &["result"], + ).unwrap(); + REGISTRY.register(Box::new(git_push_authorization.clone())).unwrap(); + + // Nostr event metrics + let events_received_total = CounterVec::new( + Opts::new( + "ngit_events_received_total", + "Nostr events received by kind", + ), + &["kind"], + ).unwrap(); + REGISTRY.register(Box::new(events_received_total.clone())).unwrap(); + + let events_stored_total = CounterVec::new( + Opts::new( + "ngit_events_stored_total", + "Nostr events successfully stored by kind", + ), + &["kind"], + ).unwrap(); + REGISTRY.register(Box::new(events_stored_total.clone())).unwrap(); + + let events_rejected_total = CounterVec::new( + Opts::new( + "ngit_events_rejected_total", + "Nostr events rejected by kind and reason", + ), + &["kind", "reason"], + ).unwrap(); + REGISTRY.register(Box::new(events_rejected_total.clone())).unwrap(); + + // Repository metrics + let repositories_total = Gauge::with_opts( + Opts::new( + "ngit_repositories_total", + "Total repositories hosted", + ) + ).unwrap(); + REGISTRY.register(Box::new(repositories_total.clone())).unwrap(); + + // Build info + let build_info = GaugeVec::new( + Opts::new( + "ngit_build_info", + "Build information", + ), + &["version", "commit"], + ).unwrap(); + REGISTRY.register(Box::new(build_info.clone())).unwrap(); + + // Set build info gauge to 1 (it's just for labels) + build_info + .with_label_values(&[env!("CARGO_PKG_VERSION"), option_env!("GIT_HASH").unwrap_or("unknown")]) + .set(1.0); + + Self { + connection_tracker, + bandwidth_tracker, + websocket_connections_total, + websocket_connection_duration, + websocket_messages_received, + websocket_messages_sent, + git_operations_total, + git_operation_duration, + git_bytes_total, + git_push_authorization, + events_received_total, + events_stored_total, + events_rejected_total, + repositories_total, + start_time: Instant::now(), + build_info, + } + } +} + +/// Timer for tracking WebSocket connection duration. +/// Records the elapsed time when dropped. +pub struct HistogramTimer { + histogram: Histogram, + start: Instant, +} + +impl HistogramTimer { + fn new(histogram: Histogram) -> Self { + Self { + histogram, + start: Instant::now(), + } + } +} + +impl Drop for HistogramTimer { + fn drop(&mut self) { + let elapsed = self.start.elapsed().as_secs_f64(); + self.histogram.observe(elapsed); + } +} + +/// Timer for tracking Git operation duration. +/// Records the elapsed time when dropped. +pub struct GitOperationTimer { + histogram_vec: HistogramVec, + operation: String, + start: Instant, +} + +impl GitOperationTimer { + fn new(histogram_vec: HistogramVec, operation: String) -> Self { + Self { + histogram_vec, + operation, + start: Instant::now(), + } + } +} + +impl Drop for GitOperationTimer { + fn drop(&mut self) { + let elapsed = self.start.elapsed().as_secs_f64(); + self.histogram_vec + .with_label_values(&[&self.operation]) + .observe(elapsed); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_metrics_creation() { + // Note: This test may fail if run with other tests due to global registry + // In production, consider using a test-specific registry + let metrics = Metrics::new(10); + + // Test that we can record metrics without panicking + metrics.record_websocket_connection(); + metrics.record_message_received("REQ"); + metrics.record_message_sent("EVENT"); + metrics.record_git_operation("clone", "success"); + metrics.record_git_bytes("in", 1024); + metrics.record_event_received(1); + metrics.record_event_stored(1); + metrics.record_event_rejected(1, "invalid_signature"); + metrics.set_repositories_total(5); + } +} \ No newline at end of file -- cgit v1.2.3