diff --git a/NOISE_FLOOR_USAGE.md b/NOISE_FLOOR_USAGE.md deleted file mode 100644 index 3934c31..0000000 --- a/NOISE_FLOOR_USAGE.md +++ /dev/null @@ -1,153 +0,0 @@ -# Noise Floor Measurement - Usage Guide - -## Overview -The noise floor measurement capability has been added to provide real-time RF environment monitoring. - -## Implementation - -### 1. Radio Wrapper (pyMC_core) -Added `get_noise_floor()` method to `SX1262Radio` class: - -```python -def get_noise_floor(self) -> Optional[float]: - """ - Get current noise floor (instantaneous RSSI) in dBm. - Returns None if radio is not initialized or if reading fails. - """ -``` - -### 2. Repeater Engine (pyMC_Repeater) -Added `get_noise_floor()` method to `RepeaterHandler` class: - -```python -def get_noise_floor(self) -> Optional[float]: - """ - Get the current noise floor (instantaneous RSSI) from the radio in dBm. - Returns None if radio is not available or reading fails. - """ -``` - -The noise floor is automatically included in the stats dictionary returned by `get_stats()`: - -```python -stats = handler.get_stats() -noise_floor = stats.get('noise_floor_dbm') # Returns float or None -``` - -## Usage Examples - -### Example 1: Get Noise Floor Directly -```python -# From the repeater engine -handler = RepeaterHandler(config, dispatcher, local_hash) -noise_floor_dbm = handler.get_noise_floor() - -if noise_floor_dbm is not None: - print(f"Current noise floor: {noise_floor_dbm:.1f} dBm") -else: - print("Noise floor not available") -``` - -### Example 2: Access via Stats -```python -# Get all stats including noise floor -stats = handler.get_stats() -noise_floor = stats.get('noise_floor_dbm') - -if noise_floor is not None: - print(f"RF Environment: {noise_floor:.1f} dBm") -``` - -### Example 3: Monitor RF Environment -```python -import asyncio - -async def monitor_rf_environment(handler, interval=5.0): - """Monitor noise floor every N seconds""" - while True: - noise_floor = handler.get_noise_floor() - if noise_floor is not None: - if noise_floor > -100: - print(f"⚠️ High RF noise: {noise_floor:.1f} dBm") - else: - print(f"✓ Normal RF environment: {noise_floor:.1f} dBm") - await asyncio.sleep(interval) -``` - -### Example 4: Channel Assessment Before TX -```python -async def should_transmit(handler, threshold_dbm=-110): - """ - Check if channel is clear before transmitting. - Returns True if noise floor is below threshold (channel clear). - """ - noise_floor = handler.get_noise_floor() - - if noise_floor is None: - # Can't determine, allow transmission - return True - - if noise_floor > threshold_dbm: - # Channel busy - high noise - print(f"Channel busy: {noise_floor:.1f} dBm > {threshold_dbm} dBm") - return False - - # Channel clear - return True -``` - -## Integration with Web Dashboard - -The noise floor is automatically available in the `/api/stats` endpoint: - -```javascript -// JavaScript example for web dashboard -fetch('/api/stats') - .then(response => response.json()) - .then(data => { - const noiseFloor = data.noise_floor_dbm; - if (noiseFloor !== null) { - updateNoiseFloorDisplay(noiseFloor); - } - }); -``` - -## Interpretation - -### Typical Values -- **-120 to -110 dBm**: Very quiet RF environment (rural, low interference) -- **-110 to -100 dBm**: Normal RF environment (typical conditions) -- **-100 to -90 dBm**: Moderate RF noise (urban, some interference) -- **-90 dBm and above**: High RF noise (congested environment, potential issues) - -### Use Cases -1. **Collision Avoidance**: Check noise floor before transmitting to detect if another station is already transmitting -2. **RF Environment Monitoring**: Track RF noise levels over time for site assessment -3. **Adaptive Transmission**: Adjust TX timing or power based on channel conditions -4. **Debugging**: Identify sources of interference or poor reception - -## Technical Details - -### Calculation -The noise floor is calculated from the SX1262's instantaneous RSSI register: -```python -raw_rssi = self.lora.getRssiInst() -noise_floor_dbm = -(float(raw_rssi) / 2) -``` - -### Update Rate -The noise floor is read on-demand when `get_noise_floor()` is called. There is no caching - each call queries the radio hardware directly. - -### Error Handling -- Returns `None` if radio is not initialized -- Returns `None` if read fails (hardware error) -- Logs debug message on error (doesn't raise exceptions) - -## Future Enhancements - -Potential future improvements: -1. **Averaging**: Average noise floor over multiple samples for stability -2. **History**: Track noise floor history for trend analysis -3. **Thresholds**: Configurable thresholds for channel busy detection -4. **Carrier Sense**: Automatic carrier sense before each transmission -5. **Spectral Analysis**: Extended to include RSSI across multiple channels diff --git a/config.yaml.example b/config.yaml.example index 5b7288e..a594677 100644 --- a/config.yaml.example +++ b/config.yaml.example @@ -96,6 +96,40 @@ duty_cycle: # Maximum airtime per minute in milliseconds max_airtime_per_minute: 3600 + +# Storage Configuration +storage: + # Directory for persistent storage files (SQLite, RRD) + storage_dir: "/var/lib/pymc_repeater" + + # MQTT publishing configuration (optional) + mqtt: + # Enable/disable MQTT publishing + enabled: false + + # MQTT broker settings + broker: "localhost" + port: 1883 + + # Authentication (optional) + username: null + password: null + + # Base topic for publishing + # Messages will be published to: {base_topic}/{node_name}/{packet|advert} + base_topic: "meshcore/repeater" + + # Data retention settings + retention: + # Clean up SQLite records older than this many days + sqlite_cleanup_days: 31 + + # RRD archives are managed automatically: + # - 1 minute resolution for 1 week + # - 5 minute resolution for 1 month + # - 1 hour resolution for 1 year + + logging: # Log level: DEBUG, INFO, WARNING, ERROR level: INFO diff --git a/pyproject.toml b/pyproject.toml index 5f17701..42f482a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -34,6 +34,8 @@ dependencies = [ "pymc_core[hardware]>=1.0.4", "pyyaml>=6.0.0", "cherrypy>=18.0.0", + "rrdtool>=0.1.16", + "paho-mqtt>=1.6.0", ] diff --git a/repeater/engine.py b/repeater/engine.py index a342c90..a456cf6 100644 --- a/repeater/engine.py +++ b/repeater/engine.py @@ -16,6 +16,7 @@ from pymc_core.protocol.constants import ( from pymc_core.protocol.packet_utils import PacketHeaderUtils, PacketTimingUtils from repeater.airtime import AirtimeManager +from repeater.storage import StorageCollector logger = logging.getLogger("RepeaterHandler") @@ -69,10 +70,15 @@ class RepeaterHandler(BaseHandler): self.dropped_count = 0 self.recent_packets = [] self.max_recent_packets = 50 - self.start_time = time.time() # For uptime calculation + self.start_time = time.time() - # Neighbor tracking (repeaters discovered via adverts) - self.neighbors = {} + # Storage collector for persistent packet logging + try: + self.storage = StorageCollector(config) + logger.info("StorageCollector initialized successfully") + except Exception as e: + logger.error(f"Failed to initialize StorageCollector: {e}") + self.storage = None async def __call__(self, packet: Packet, metadata: Optional[dict] = None) -> None: @@ -221,6 +227,15 @@ class RepeaterHandler(BaseHandler): } + # Store packet record to persistent storage + if self.storage: + try: + self.storage.record_packet(packet_record) + except Exception as e: + logger.error(f"Failed to store packet record: {e}") + + + # If this is a duplicate, try to attach it to the original packet if is_dupe and len(self.recent_packets) > 0: # Find the original packet with same hash @@ -295,6 +310,7 @@ class RepeaterHandler(BaseHandler): decode_appdata, get_contact_type_name, parse_advert_payload, + determine_contact_type_from_flags, ) # Parse advert payload @@ -317,14 +333,8 @@ class RepeaterHandler(BaseHandler): appdata_decoded = decode_appdata(appdata) flags = appdata_decoded.get("flags", 0) - is_repeater = bool(flags & ADVERT_FLAG_IS_REPEATER) - - if not is_repeater: - return # Not a repeater, skip - - from pymc_core.protocol.utils import determine_contact_type_from_flags - + route_type = packet.header & PH_ROUTE_MASK contact_type_id = determine_contact_type_from_flags(flags) contact_type = get_contact_type_name(contact_type_id) @@ -334,32 +344,34 @@ class RepeaterHandler(BaseHandler): longitude = appdata_decoded.get("longitude") current_time = time.time() + + # Check if this is a new neighbor + current_neighbors = self.storage.get_neighbors() if self.storage else {} + is_new_neighbor = pubkey not in current_neighbors - # Update or create neighbor entry - if pubkey not in self.neighbors: - self.neighbors[pubkey] = { - "node_name": node_name, - "contact_type": contact_type, - "latitude": latitude, - "longitude": longitude, - "first_seen": current_time, - "last_seen": current_time, - "rssi": rssi, - "snr": snr, - "advert_count": 1, - } - logger.info(f"Discovered new repeater: {node_name} ({pubkey[:16]}...)") - else: - # Update existing neighbor - neighbor = self.neighbors[pubkey] - neighbor["node_name"] = node_name # Update name in case it changed - neighbor["contact_type"] = contact_type - neighbor["latitude"] = latitude - neighbor["longitude"] = longitude - neighbor["last_seen"] = current_time - neighbor["rssi"] = rssi - neighbor["snr"] = snr - neighbor["advert_count"] = neighbor.get("advert_count", 0) + 1 + # Create advert record for storage + advert_record = { + "timestamp": current_time, + "pubkey": pubkey, + "node_name": node_name, + "is_repeater": is_repeater, + "route_type": route_type, + "contact_type": contact_type, + "latitude": latitude, + "longitude": longitude, + "rssi": rssi, + "snr": snr, + "is_new_neighbor": is_new_neighbor, + } + + # Store to database + if self.storage: + try: + self.storage.record_advert(advert_record) + if is_new_neighbor: + logger.info(f"Discovered new neighbor: {node_name} ({pubkey[:16]}...)") + except Exception as e: + logger.error(f"Failed to store advert record: {e}") except Exception as e: logger.debug(f"Error processing advert for neighbor tracking: {e}") @@ -610,6 +622,9 @@ class RepeaterHandler(BaseHandler): # Get current noise floor from radio noise_floor_dbm = self.get_noise_floor() + # Get neighbors from database + neighbors = self.storage.get_neighbors() if self.storage else {} + stats = { "local_hash": f"0x{self.local_hash:02x}", "duplicate_cache_size": len(self.seen_packets), @@ -620,7 +635,7 @@ class RepeaterHandler(BaseHandler): "rx_per_hour": rx_per_hour, "forwarded_per_hour": forwarded_per_hour, "recent_packets": self.recent_packets, - "neighbors": self.neighbors, + "neighbors": neighbors, "uptime_seconds": uptime_seconds, "noise_floor_dbm": noise_floor_dbm, # Add configuration data @@ -656,3 +671,17 @@ class RepeaterHandler(BaseHandler): # Add airtime stats stats.update(self.airtime_mgr.get_stats()) return stats + + def cleanup(self): + if self.storage: + try: + self.storage.close() + logger.info("StorageCollector closed successfully") + except Exception as e: + logger.error(f"Error closing StorageCollector: {e}") + + def __del__(self): + try: + self.cleanup() + except Exception: + pass \ No newline at end of file diff --git a/repeater/http/__init__.py b/repeater/http/__init__.py new file mode 100644 index 0000000..77ae3f9 --- /dev/null +++ b/repeater/http/__init__.py @@ -0,0 +1,12 @@ +from .http_server import HTTPStatsServer, StatsApp, LogBuffer, _log_buffer +from .api_endpoints import APIEndpoints +from .cad_calibration_engine import CADCalibrationEngine + +__all__ = [ + 'HTTPStatsServer', + 'StatsApp', + 'LogBuffer', + 'APIEndpoints', + 'CADCalibrationEngine', + '_log_buffer' +] \ No newline at end of file diff --git a/repeater/http/api_endpoints.py b/repeater/http/api_endpoints.py new file mode 100644 index 0000000..1a76e91 --- /dev/null +++ b/repeater/http/api_endpoints.py @@ -0,0 +1,502 @@ +import json +import logging +import time +from datetime import datetime +from typing import Callable, Optional +import cherrypy +from repeater import __version__ +from .cad_calibration_engine import CADCalibrationEngine + +logger = logging.getLogger("HTTPServer") + +class APIEndpoints: + def __init__(self, stats_getter: Optional[Callable] = None, send_advert_func: Optional[Callable] = None, config: Optional[dict] = None, event_loop=None, daemon_instance=None, config_path=None): + self.stats_getter = stats_getter + self.send_advert_func = send_advert_func + self.config = config or {} + self.event_loop = event_loop + self.daemon_instance = daemon_instance + self._config_path = config_path or '/etc/pymc_repeater/config.yaml' + self.cad_calibration = CADCalibrationEngine(daemon_instance, event_loop) + + def _get_storage(self): + if not self.daemon_instance or not hasattr(self.daemon_instance, 'storage'): + raise Exception("Storage not available") + return self.daemon_instance.storage + + def _success(self, data, **kwargs): + result = {"success": True, "data": data} + result.update(kwargs) + return result + + def _error(self, error): + return {"success": False, "error": str(error)} + + def _get_params(self, defaults): + params = cherrypy.request.params + result = {} + for key, default in defaults.items(): + value = params.get(key, default) + if isinstance(default, int): + result[key] = int(value) if value is not None else None + elif isinstance(default, float): + result[key] = float(value) if value is not None else None + else: + result[key] = value + return result + + def _require_post(self): + if cherrypy.request.method != "POST": + raise Exception("Method not allowed") + + def _get_time_range(self, hours): + end_time = int(time.time()) + return end_time - (hours * 3600), end_time + + def _process_counter_data(self, data_points, timestamps_ms): + rates = [] + prev_value = None + for value in data_points: + if value is None: + rates.append(0) + elif prev_value is None: + rates.append(0) + else: + rates.append(max(0, value - prev_value)) + prev_value = value + return [[timestamps_ms[i], rates[i]] for i in range(min(len(rates), len(timestamps_ms)))] + + def _process_gauge_data(self, data_points, timestamps_ms): + values = [v if v is not None else 0 for v in data_points] + return [[timestamps_ms[i], values[i]] for i in range(min(len(values), len(timestamps_ms)))] + + @cherrypy.expose + @cherrypy.tools.json_out() + def stats(self): + try: + stats = self.stats_getter() if self.stats_getter else {} + stats["version"] = __version__ + try: + import pymc_core + stats["core_version"] = pymc_core.__version__ + except ImportError: + stats["core_version"] = "unknown" + return stats + except Exception as e: + logger.error(f"Error serving stats: {e}") + return {"error": str(e)} + + @cherrypy.expose + @cherrypy.tools.json_out() + def send_advert(self): + try: + self._require_post() + if not self.send_advert_func: + return self._error("Send advert function not configured") + if self.event_loop is None: + return self._error("Event loop not available") + import asyncio + future = asyncio.run_coroutine_threadsafe(self.send_advert_func(), self.event_loop) + result = future.result(timeout=10) + return self._success("Advert sent successfully") if result else self._error("Failed to send advert") + except Exception as e: + logger.error(f"Error sending advert: {e}", exc_info=True) + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + @cherrypy.tools.json_in() + def set_mode(self): + try: + self._require_post() + data = cherrypy.request.json + new_mode = data.get("mode", "forward") + if new_mode not in ["forward", "monitor"]: + return self._error("Invalid mode. Must be 'forward' or 'monitor'") + if "repeater" not in self.config: + self.config["repeater"] = {} + self.config["repeater"]["mode"] = new_mode + logger.info(f"Mode changed to: {new_mode}") + return {"success": True, "mode": new_mode} + except Exception as e: + logger.error(f"Error setting mode: {e}", exc_info=True) + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + @cherrypy.tools.json_in() + def set_duty_cycle(self): + try: + self._require_post() + data = cherrypy.request.json + enabled = data.get("enabled", True) + if "duty_cycle" not in self.config: + self.config["duty_cycle"] = {} + self.config["duty_cycle"]["enforcement_enabled"] = enabled + logger.info(f"Duty cycle enforcement {'enabled' if enabled else 'disabled'}") + return {"success": True, "enabled": enabled} + except Exception as e: + logger.error(f"Error setting duty cycle: {e}", exc_info=True) + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def logs(self): + from .http_server import _log_buffer + try: + logs = list(_log_buffer.logs) + return { + "logs": ( + logs + if logs + else [ + { + "message": "No logs available", + "timestamp": datetime.now().isoformat(), + "level": "INFO", + } + ] + ) + } + except Exception as e: + logger.error(f"Error fetching logs: {e}") + return {"error": str(e), "logs": []} + + @cherrypy.expose + @cherrypy.tools.json_out() + def packet_stats(self): + try: + hours = int(cherrypy.request.params.get('hours', 24)) + stats = self._get_storage().get_packet_stats(hours=hours) + return self._success(stats) + except Exception as e: + logger.error(f"Error getting packet stats: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def recent_packets(self): + try: + limit = int(cherrypy.request.params.get('limit', 100)) + packets = self._get_storage().get_recent_packets(limit=limit) + return self._success(packets, count=len(packets)) + except Exception as e: + logger.error(f"Error getting recent packets: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def filtered_packets(self): + try: + params = self._get_params({ + 'type': None, + 'route': None, + 'start_timestamp': None, + 'end_timestamp': None, + 'limit': 1000 + }) + packets = self._get_storage().get_filtered_packets(**params) + return self._success(packets, count=len(packets), filters=params) + except ValueError as e: + return self._error(f"Invalid parameter format: {e}") + except Exception as e: + logger.error(f"Error getting filtered packets: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def packet_by_hash(self, packet_hash=None): + try: + if not packet_hash: + return self._error("packet_hash parameter required") + packet = self._get_storage().get_packet_by_hash(packet_hash) + return self._success(packet) if packet else self._error("Packet not found") + except Exception as e: + logger.error(f"Error getting packet by hash: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def packet_type_stats(self): + try: + hours = int(cherrypy.request.params.get('hours', 24)) + stats = self._get_storage().get_packet_type_stats(hours=hours) + return self._success(stats) + except Exception as e: + logger.error(f"Error getting packet type stats: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def rrd_data(self): + try: + params = self._get_params({ + 'start_time': None, + 'end_time': None, + 'resolution': 'average' + }) + data = self._get_storage().get_rrd_data(**params) + return self._success(data) if data else self._error("No RRD data available") + except ValueError as e: + return self._error(f"Invalid parameter format: {e}") + except Exception as e: + logger.error(f"Error getting RRD data: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def packet_type_graph_data(self): + try: + params = self._get_params({'hours': 24, 'resolution': 'average', 'types': 'all'}) + start_time, end_time = self._get_time_range(params['hours']) + + rrd_data = self._get_storage().get_rrd_data( + start_time=start_time, end_time=end_time, resolution=params['resolution'] + ) + + if not rrd_data or 'packet_types' not in rrd_data: + return self._error("No RRD data available") + + packet_type_names = { + 'type_0': 'Request (REQ)', 'type_1': 'Response (RESPONSE)', + 'type_2': 'Text Message (TXT_MSG)', 'type_3': 'ACK (ACK)', + 'type_4': 'Advert (ADVERT)', 'type_5': 'Group Text (GRP_TXT)', + 'type_6': 'Group Data (GRP_DATA)', 'type_7': 'Anonymous Request (ANON_REQ)', + 'type_8': 'Path (PATH)', 'type_9': 'Trace (TRACE)', + 'type_10': 'Reserved Type 10', 'type_11': 'Reserved Type 11', + 'type_12': 'Reserved Type 12', 'type_13': 'Reserved Type 13', + 'type_14': 'Reserved Type 14', 'type_15': 'Reserved Type 15', + 'type_other': 'Other Types (>15)' + } + + if params['types'] != 'all': + requested_types = [f'type_{t.strip()}' for t in params['types'].split(',')] + if 'other' in params['types'].lower(): + requested_types.append('type_other') + else: + requested_types = list(rrd_data['packet_types'].keys()) + + timestamps_ms = [ts * 1000 for ts in rrd_data['timestamps']] + series = [] + + for type_key in requested_types: + if type_key in rrd_data['packet_types']: + chart_data = self._process_counter_data(rrd_data['packet_types'][type_key], timestamps_ms) + series.append({ + "name": packet_type_names.get(type_key, type_key), + "type": type_key, + "data": chart_data + }) + + graph_data = { + "start_time": rrd_data['start_time'], + "end_time": rrd_data['end_time'], + "step": rrd_data['step'], + "timestamps": rrd_data['timestamps'], + "series": series + } + + return self._success(graph_data) + + except ValueError as e: + return self._error(f"Invalid parameter format: {e}") + except Exception as e: + logger.error(f"Error getting packet type graph data: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def metrics_graph_data(self): + try: + params = self._get_params({'hours': 24, 'resolution': 'average', 'metrics': 'all'}) + start_time, end_time = self._get_time_range(params['hours']) + + rrd_data = self._get_storage().get_rrd_data( + start_time=start_time, end_time=end_time, resolution=params['resolution'] + ) + + if not rrd_data or 'metrics' not in rrd_data: + return self._error("No RRD data available") + + metric_names = { + 'rx_count': 'Received Packets', 'tx_count': 'Transmitted Packets', + 'drop_count': 'Dropped Packets', 'avg_rssi': 'Average RSSI (dBm)', + 'avg_snr': 'Average SNR (dB)', 'avg_length': 'Average Packet Length', + 'avg_score': 'Average Score', 'neighbor_count': 'Neighbor Count' + } + + counter_metrics = ['rx_count', 'tx_count', 'drop_count'] + + if params['metrics'] != 'all': + requested_metrics = [m.strip() for m in params['metrics'].split(',')] + else: + requested_metrics = list(rrd_data['metrics'].keys()) + + timestamps_ms = [ts * 1000 for ts in rrd_data['timestamps']] + series = [] + + for metric_key in requested_metrics: + if metric_key in rrd_data['metrics']: + if metric_key in counter_metrics: + chart_data = self._process_counter_data(rrd_data['metrics'][metric_key], timestamps_ms) + else: + chart_data = self._process_gauge_data(rrd_data['metrics'][metric_key], timestamps_ms) + + series.append({ + "name": metric_names.get(metric_key, metric_key), + "type": metric_key, + "data": chart_data + }) + + graph_data = { + "start_time": rrd_data['start_time'], + "end_time": rrd_data['end_time'], + "step": rrd_data['step'], + "timestamps": rrd_data['timestamps'], + "series": series + } + + return self._success(graph_data) + + except ValueError as e: + return self._error(f"Invalid parameter format: {e}") + except Exception as e: + logger.error(f"Error getting metrics graph data: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + @cherrypy.tools.json_in() + def cad_calibration_start(self): + try: + self._require_post() + data = cherrypy.request.json or {} + samples = data.get("samples", 8) + delay = data.get("delay", 100) + if self.cad_calibration.start_calibration(samples, delay): + return self._success("Calibration started") + else: + return self._error("Calibration already running") + except Exception as e: + logger.error(f"Error starting CAD calibration: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + def cad_calibration_stop(self): + try: + self._require_post() + self.cad_calibration.stop_calibration() + return self._success("Calibration stopped") + except Exception as e: + logger.error(f"Error stopping CAD calibration: {e}") + return self._error(e) + + @cherrypy.expose + @cherrypy.tools.json_out() + @cherrypy.tools.json_in() + def save_cad_settings(self): + try: + self._require_post() + data = cherrypy.request.json or {} + peak = data.get("peak") + min_val = data.get("min_val") + detection_rate = data.get("detection_rate", 0) + + if peak is None or min_val is None: + return self._error("Missing peak or min_val parameters") + + if self.daemon_instance and hasattr(self.daemon_instance, 'radio') and self.daemon_instance.radio: + if hasattr(self.daemon_instance.radio, 'set_custom_cad_thresholds'): + self.daemon_instance.radio.set_custom_cad_thresholds(peak=peak, min_val=min_val) + logger.info(f"Applied CAD settings to radio: peak={peak}, min={min_val}") + + if "radio" not in self.config: + self.config["radio"] = {} + if "cad" not in self.config["radio"]: + self.config["radio"]["cad"] = {} + + self.config["radio"]["cad"]["peak_threshold"] = peak + self.config["radio"]["cad"]["min_threshold"] = min_val + + config_path = getattr(self, '_config_path', '/etc/pymc_repeater/config.yaml') + self._save_config_to_file(config_path) + + logger.info(f"Saved CAD settings to config: peak={peak}, min={min_val}, rate={detection_rate:.1f}%") + return { + "success": True, + "message": f"CAD settings saved: peak={peak}, min={min_val}", + "settings": {"peak": peak, "min_val": min_val, "detection_rate": detection_rate} + } + except Exception as e: + logger.error(f"Error saving CAD settings: {e}") + return self._error(e) + + def _save_config_to_file(self, config_path): + try: + import yaml + import os + os.makedirs(os.path.dirname(config_path), exist_ok=True) + with open(config_path, 'w') as f: + yaml.dump(self.config, f, default_flow_style=False, indent=2) + logger.info(f"Configuration saved to {config_path}") + except Exception as e: + logger.error(f"Failed to save config to {config_path}: {e}") + raise + + @cherrypy.expose + def cad_calibration_stream(self): + cherrypy.response.headers['Content-Type'] = 'text/event-stream' + cherrypy.response.headers['Cache-Control'] = 'no-cache' + cherrypy.response.headers['Connection'] = 'keep-alive' + cherrypy.response.headers['Access-Control-Allow-Origin'] = '*' + + if not hasattr(self.cad_calibration, 'message_queue'): + self.cad_calibration.message_queue = [] + + def generate(): + try: + yield f"data: {json.dumps({'type': 'connected', 'message': 'Connected to CAD calibration stream'})}\n\n" + + if self.cad_calibration.running: + config = getattr(self.cad_calibration.daemon_instance, 'config', {}) + radio_config = config.get("radio", {}) + sf = radio_config.get("spreading_factor", 8) + + peak_range, min_range = self.cad_calibration.get_test_ranges(sf) + total_tests = len(peak_range) * len(min_range) + + status_message = { + "type": "status", + "message": f"Calibration in progress: SF{sf}, {total_tests} tests", + "test_ranges": { + "peak_min": min(peak_range), + "peak_max": max(peak_range), + "min_min": min(min_range), + "min_max": max(min_range), + "spreading_factor": sf, + "total_tests": total_tests + } + } + yield f"data: {json.dumps(status_message)}\n\n" + + last_message_index = len(self.cad_calibration.message_queue) + + while True: + current_queue_length = len(self.cad_calibration.message_queue) + if current_queue_length > last_message_index: + for i in range(last_message_index, current_queue_length): + message = self.cad_calibration.message_queue[i] + yield f"data: {json.dumps(message)}\n\n" + last_message_index = current_queue_length + else: + yield f"data: {json.dumps({'type': 'keepalive'})}\n\n" + + time.sleep(0.5) + + except Exception as e: + logger.error(f"SSE stream error: {e}") + + return generate() + + cad_calibration_stream._cp_config = {'response.stream': True} \ No newline at end of file diff --git a/repeater/http/cad_calibration_engine.py b/repeater/http/cad_calibration_engine.py new file mode 100644 index 0000000..c124cfe --- /dev/null +++ b/repeater/http/cad_calibration_engine.py @@ -0,0 +1,287 @@ +import asyncio +import logging +import random +import threading +import time +from typing import Dict, Any, Optional + +logger = logging.getLogger("HTTPServer") + + +class CADCalibrationEngine: + + def __init__(self, daemon_instance=None, event_loop=None): + self.daemon_instance = daemon_instance + self.event_loop = event_loop + self.running = False + self.results = {} + self.current_test = None + self.progress = {"current": 0, "total": 0} + self.clients = set() # SSE clients + self.calibration_thread = None + + def get_test_ranges(self, spreading_factor: int): + """Get CAD test ranges""" + # Higher values = less sensitive, lower values = more sensitive + # Test from LESS sensitive to MORE sensitive to find the sweet spot + sf_ranges = { + 7: (range(22, 30, 1), range(12, 20, 1)), + 8: (range(22, 30, 1), range(12, 20, 1)), + 9: (range(24, 32, 1), range(14, 22, 1)), + 10: (range(26, 34, 1), range(16, 24, 1)), + 11: (range(28, 36, 1), range(18, 26, 1)), + 12: (range(30, 38, 1), range(20, 28, 1)), + } + return sf_ranges.get(spreading_factor, sf_ranges[8]) + + async def test_cad_config(self, radio, det_peak: int, det_min: int, samples: int = 20) -> Dict[str, Any]: + + detections = 0 + baseline_detections = 0 + + # First, get baseline with very insensitive settings (should detect nothing) + baseline_samples = 5 + for _ in range(baseline_samples): + try: + # Use very high thresholds that should detect nothing + baseline_result = await radio.perform_cad(det_peak=35, det_min=25, timeout=0.3) + if baseline_result: + baseline_detections += 1 + except Exception: + pass + await asyncio.sleep(0.1) # 100ms between baseline samples + + # Wait before actual test + await asyncio.sleep(0.5) + + # Now test the actual configuration + for i in range(samples): + try: + result = await radio.perform_cad(det_peak=det_peak, det_min=det_min, timeout=0.3) + if result: + detections += 1 + except Exception: + pass + + # Variable delay to avoid sampling artifacts + delay = 0.05 + (i % 3) * 0.05 # 50ms, 100ms, 150ms rotation + await asyncio.sleep(delay) + + # Calculate adjusted detection rate + baseline_rate = (baseline_detections / baseline_samples) * 100 + detection_rate = (detections / samples) * 100 + + # Subtract baseline noise + adjusted_rate = max(0, detection_rate - baseline_rate) + + return { + 'det_peak': det_peak, + 'det_min': det_min, + 'samples': samples, + 'detections': detections, + 'detection_rate': detection_rate, + 'baseline_rate': baseline_rate, + 'adjusted_rate': adjusted_rate, # This is the useful metric + 'sensitivity_score': self._calculate_sensitivity_score(det_peak, det_min, adjusted_rate) + } + + def _calculate_sensitivity_score(self, det_peak: int, det_min: int, adjusted_rate: float) -> float: + + # Ideal detection rate is around 10-30% for good sensitivity without false positives + ideal_rate = 20.0 + rate_penalty = abs(adjusted_rate - ideal_rate) / ideal_rate + + # Prefer moderate sensitivity settings (not too extreme) + sensitivity_penalty = (abs(det_peak - 25) + abs(det_min - 15)) / 20.0 + + # Lower penalty = higher score + score = max(0, 100 - (rate_penalty * 50) - (sensitivity_penalty * 20)) + return score + + def broadcast_to_clients(self, data): + + # Store the message for clients to pick up + self.last_message = data + # Also store in a queue for clients to consume + if not hasattr(self, 'message_queue'): + self.message_queue = [] + self.message_queue.append(data) + + def calibration_worker(self, samples: int, delay_ms: int): + + try: + # Get radio from daemon instance + if not self.daemon_instance: + self.broadcast_to_clients({"type": "error", "message": "No daemon instance available"}) + return + + radio = getattr(self.daemon_instance, 'radio', None) + if not radio: + self.broadcast_to_clients({"type": "error", "message": "Radio instance not available"}) + return + if not hasattr(radio, 'perform_cad'): + self.broadcast_to_clients({"type": "error", "message": "Radio does not support CAD"}) + return + + # Get spreading factor from daemon instance + config = getattr(self.daemon_instance, 'config', {}) + radio_config = config.get("radio", {}) + sf = radio_config.get("spreading_factor", 8) + + # Get test ranges + peak_range, min_range = self.get_test_ranges(sf) + + total_tests = len(peak_range) * len(min_range) + self.progress = {"current": 0, "total": total_tests} + + self.broadcast_to_clients({ + "type": "status", + "message": f"Starting calibration: SF{sf}, {total_tests} tests", + "test_ranges": { + "peak_min": min(peak_range), + "peak_max": max(peak_range), + "min_min": min(min_range), + "min_max": max(min_range), + "spreading_factor": sf, + "total_tests": total_tests + } + }) + + current = 0 + + peak_list = list(peak_range) + min_list = list(min_range) + + # Create all test combinations + test_combinations = [] + for det_peak in peak_list: + for det_min in min_list: + test_combinations.append((det_peak, det_min)) + + # Sort by distance from center for center-out pattern + peak_center = (max(peak_list) + min(peak_list)) / 2 + min_center = (max(min_list) + min(min_list)) / 2 + + def distance_from_center(combo): + peak, min_val = combo + return ((peak - peak_center) ** 2 + (min_val - min_center) ** 2) ** 0.5 + + # Sort by distance from center + test_combinations.sort(key=distance_from_center) + + # Randomize within bands for better coverage + band_size = max(1, len(test_combinations) // 8) # Create 8 bands + randomized_combinations = [] + + for i in range(0, len(test_combinations), band_size): + band = test_combinations[i:i + band_size] + random.shuffle(band) # Randomize within each band + randomized_combinations.extend(band) + + # Run calibration in event loop with center-out randomized pattern + if self.event_loop: + for det_peak, det_min in randomized_combinations: + if not self.running: + break + + current += 1 + self.progress["current"] = current + + # Update progress + self.broadcast_to_clients({ + "type": "progress", + "current": current, + "total": total_tests, + "peak": det_peak, + "min": det_min + }) + + # Run the test + future = asyncio.run_coroutine_threadsafe( + self.test_cad_config(radio, det_peak, det_min, samples), + self.event_loop + ) + + try: + result = future.result(timeout=30) # 30 second timeout per test + + # Store result + key = f"{det_peak}-{det_min}" + self.results[key] = result + + # Send result to clients + self.broadcast_to_clients({ + "type": "result", + **result + }) + except Exception as e: + logger.error(f"CAD test failed for peak={det_peak}, min={det_min}: {e}") + + # Delay between tests + if self.running and delay_ms > 0: + time.sleep(delay_ms / 1000.0) + + if self.running: + # Find best result based on sensitivity score (not just detection rate) + best_result = None + recommended_result = None + if self.results: + # Find result with highest sensitivity score (best balance) + best_result = max(self.results.values(), key=lambda x: x.get('sensitivity_score', 0)) + + # Also find result with ideal adjusted detection rate (10-30%) + ideal_results = [r for r in self.results.values() if 10 <= r.get('adjusted_rate', 0) <= 30] + if ideal_results: + # Among ideal results, pick the one with best sensitivity score + recommended_result = max(ideal_results, key=lambda x: x.get('sensitivity_score', 0)) + else: + recommended_result = best_result + + self.broadcast_to_clients({ + "type": "completed", + "message": "Calibration completed", + "results": { + "best": best_result, + "recommended": recommended_result, + "total_tests": len(self.results) + } if best_result else None + }) + else: + self.broadcast_to_clients({"type": "status", "message": "Calibration stopped"}) + + except Exception as e: + logger.error(f"Calibration worker error: {e}") + self.broadcast_to_clients({"type": "error", "message": str(e)}) + finally: + self.running = False + + def start_calibration(self, samples: int = 8, delay_ms: int = 100): + + if self.running: + return False + + self.running = True + self.results.clear() + self.progress = {"current": 0, "total": 0} + self.clear_message_queue() # Clear any old messages + + # Start calibration in separate thread + self.calibration_thread = threading.Thread( + target=self.calibration_worker, + args=(samples, delay_ms) + ) + self.calibration_thread.daemon = True + self.calibration_thread.start() + + return True + + def stop_calibration(self): + + self.running = False + if self.calibration_thread: + self.calibration_thread.join(timeout=2) + + def clear_message_queue(self): + + if hasattr(self, 'message_queue'): + self.message_queue.clear() \ No newline at end of file diff --git a/repeater/http/http_server.py b/repeater/http/http_server.py new file mode 100644 index 0000000..8457db6 --- /dev/null +++ b/repeater/http/http_server.py @@ -0,0 +1,342 @@ +import json +import logging +import os +import re +from collections import deque +from datetime import datetime +from typing import Callable, Optional + +import cherrypy +from pymc_core.protocol.utils import PAYLOAD_TYPES, ROUTE_TYPES + +from repeater import __version__ +from .api_endpoints import APIEndpoints + +logger = logging.getLogger("HTTPServer") + + +# In-memory log buffer +class LogBuffer(logging.Handler): + + def __init__(self, max_lines=100): + super().__init__() + self.logs = deque(maxlen=max_lines) + self.setFormatter(logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")) + + def emit(self, record): + + try: + msg = self.format(record) + self.logs.append( + { + "message": msg, + "timestamp": datetime.fromtimestamp(record.created).isoformat(), + "level": record.levelname, + } + ) + except Exception: + self.handleError(record) + + +# Global log buffer instance +_log_buffer = LogBuffer(max_lines=100) + +class StatsApp: + + def __init__( + self, + stats_getter: Optional[Callable] = None, + template_dir: Optional[str] = None, + node_name: str = "Repeater", + pub_key: str = "", + send_advert_func: Optional[Callable] = None, + config: Optional[dict] = None, + event_loop=None, + daemon_instance=None, + config_path=None, + ): + + self.stats_getter = stats_getter + self.template_dir = template_dir + self.node_name = node_name + self.pub_key = pub_key + self.dashboard_template = None + self.config = config or {} + + # Create nested API object for routing + self.api = APIEndpoints(stats_getter, send_advert_func, self.config, event_loop, daemon_instance, config_path) + + # Load template on init + if template_dir: + template_path = os.path.join(template_dir, "dashboard.html") + try: + with open(template_path, "r") as f: + self.dashboard_template = f.read() + logger.info(f"Loaded template from {template_path}") + except FileNotFoundError: + logger.error(f"Template not found: {template_path}") + + @cherrypy.expose + def index(self): + """Serve dashboard HTML.""" + return self._serve_template("dashboard.html") + + @cherrypy.expose + def neighbors(self): + """Serve neighbors page.""" + return self._serve_template("neighbors.html") + + @cherrypy.expose + def statistics(self): + """Serve statistics page.""" + return self._serve_template("statistics.html") + + @cherrypy.expose + def configuration(self): + """Serve configuration page.""" + return self._serve_template("configuration.html") + + @cherrypy.expose + def logs(self): + """Serve logs page.""" + return self._serve_template("logs.html") + + @cherrypy.expose + def help(self): + """Serve help documentation.""" + return self._serve_template("help.html") + + @cherrypy.expose + def cad_calibration(self): + """Serve CAD calibration page.""" + return self._serve_template("cad-calibration.html") + + def _serve_template(self, template_name: str): + """Serve HTML template with stats.""" + if not self.template_dir: + return "

Error

Template directory not configured

" + + if not self.dashboard_template: + return "

Error

Template not loaded

" + + try: + + template_path = os.path.join(self.template_dir, template_name) + with open(template_path, "r") as f: + template_content = f.read() + + nav_path = os.path.join(self.template_dir, "nav.html") + nav_content = "" + try: + with open(nav_path, "r") as f: + nav_content = f.read() + except FileNotFoundError: + logger.warning(f"Navigation template not found: {nav_path}") + + stats = self.stats_getter() if self.stats_getter else {} + + if "uptime_seconds" not in stats or not isinstance( + stats.get("uptime_seconds"), (int, float) + ): + stats["uptime_seconds"] = 0 + + # Calculate uptime in hours + uptime_seconds = stats.get("uptime_seconds", 0) + uptime_hours = int(uptime_seconds // 3600) if uptime_seconds else 0 + + # Determine current page for nav highlighting + page_map = { + "dashboard.html": "dashboard", + "neighbors.html": "neighbors", + "statistics.html": "statistics", + "configuration.html": "configuration", + "cad-calibration.html": "cad-calibration", + "logs.html": "logs", + "help.html": "help", + } + current_page = page_map.get(template_name, "") + + # Prepare basic substitutions + html = template_content + html = html.replace("{{ node_name }}", str(self.node_name)) + html = html.replace("{{ last_updated }}", datetime.now().strftime("%Y-%m-%d %H:%M:%S")) + html = html.replace("{{ page }}", current_page) + + # Replace navigation placeholder with actual nav content + if "" in html: + nav_substitutions = nav_content + nav_substitutions = nav_substitutions.replace( + "{{ node_name }}", str(self.node_name) + ) + nav_substitutions = nav_substitutions.replace("{{ pub_key }}", str(self.pub_key)) + nav_substitutions = nav_substitutions.replace( + "{{ last_updated }}", datetime.now().strftime("%Y-%m-%d %H:%M:%S") + ) + + # Handle active state for nav items + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'dashboard' else '' }}", + " active" if current_page == "dashboard" else "", + ) + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'neighbors' else '' }}", + " active" if current_page == "neighbors" else "", + ) + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'statistics' else '' }}", + " active" if current_page == "statistics" else "", + ) + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'configuration' else '' }}", + " active" if current_page == "configuration" else "", + ) + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'logs' else '' }}", + " active" if current_page == "logs" else "", + ) + nav_substitutions = nav_substitutions.replace( + "{{ ' active' if page == 'help' else '' }}", + " active" if current_page == "help" else "", + ) + + html = html.replace("", nav_substitutions) + + # Build packets table HTML for dashboard + if template_name == "dashboard.html": + recent_packets = stats.get("recent_packets", []) + packets_table = "" + + if recent_packets: + for pkt in recent_packets[-20:]: # Last 20 packets + time_obj = datetime.fromtimestamp(pkt.get("timestamp", 0)) + time_str = time_obj.strftime("%H:%M:%S") + pkt_type = PAYLOAD_TYPES.get( + pkt.get("type", 0), f"0x{pkt.get('type', 0): 02x}" + ) + route_type = pkt.get("route", 0) + route = ROUTE_TYPES.get(route_type, f"UNKNOWN_{route_type}") + status = "OK TX" if pkt.get("transmitted") else "WAIT" + + # Get proper CSS class for route type + route_class = route.lower().replace("_", "-") + snr_val = pkt.get("snr", 0.0) + score_val = pkt.get("score", 0) + delay_val = pkt.get("tx_delay_ms", 0) + + packets_table += ( + "" + f"{time_str}" + f'{pkt_type}' + f'{route}' + f"{pkt.get('length', 0)}" + f"{pkt.get('rssi', 0)}" + f"{snr_val: .1f}" + f'{score_val: .2f}' + f"{delay_val: .0f}" + f"{status}" + "" + ) + else: + packets_table = """ + + + No packets received yet - waiting for traffic... + + + """ + + # Add dashboard-specific substitutions + html = html.replace("{{ rx_count }}", str(stats.get("rx_count", 0))) + html = html.replace("{{ forwarded_count }}", str(stats.get("forwarded_count", 0))) + html = html.replace("{{ dropped_count }}", str(stats.get("dropped_count", 0))) + html = html.replace("{{ uptime_hours }}", str(uptime_hours)) + + # Replace tbody with actual packets + tbody_pattern = r'.*?' + tbody_replacement = f'\n{packets_table}\n' + html = re.sub( + tbody_pattern, + tbody_replacement, + html, + flags=re.DOTALL, + ) + + return html + + except Exception as e: + logger.error(f"Error rendering template {template_name}: {e}", exc_info=True) + return f"

Error

{str(e)}

" + + +class HTTPStatsServer: + + def __init__( + self, + host: str = "0.0.0.0", + port: int = 8000, + stats_getter: Optional[Callable] = None, + template_dir: Optional[str] = None, + node_name: str = "Repeater", + pub_key: str = "", + send_advert_func: Optional[Callable] = None, + config: Optional[dict] = None, + event_loop=None, + daemon_instance=None, + config_path=None, + ): + + self.host = host + self.port = port + self.app = StatsApp( + stats_getter, template_dir, node_name, pub_key, send_advert_func, config, event_loop, daemon_instance, config_path + ) + + def start(self): + + try: + # Serve static files from templates directory + static_dir = ( + self.app.template_dir if self.app.template_dir else os.path.dirname(__file__) + ) + + config = { + "/": { + "tools.sessions.on": False, + }, + "/static": { + "tools.staticdir.on": True, + "tools.staticdir.dir": static_dir, + }, + } + + cherrypy.config.update( + { + "server.socket_host": self.host, + "server.socket_port": self.port, + "engine.autoreload.on": False, + "log.screen": False, + "log.access_file": "", # Disable access log file + "log.error_file": "", # Disable error log file + } + ) + + cherrypy.tree.mount(self.app, "/", config) + + # Completely disable access logging + cherrypy.log.access_log.propagate = False + cherrypy.log.error_log.setLevel(logging.ERROR) + + cherrypy.engine.start() + server_url = "http://{}:{}".format(self.host, self.port) + logger.info(f"HTTP stats server started on {server_url}") + + except Exception as e: + logger.error(f"Failed to start HTTP server: {e}") + raise + + def stop(self): + try: + cherrypy.engine.exit() + logger.info("HTTP stats server stopped") + except Exception as e: + logger.warning(f"Error stopping HTTP server: {e}") diff --git a/repeater/http/stats-data-collection.md b/repeater/http/stats-data-collection.md new file mode 100644 index 0000000..d1a646e --- /dev/null +++ b/repeater/http/stats-data-collection.md @@ -0,0 +1,1508 @@ +# Stats Data Collection & Charting Examples + +This document provides examples for using the pyMC_Repeater API endpoints to create charts and visualizations for network monitoring. + +## Available API Endpoints + +### Basic Statistics +- `/api/packet_stats` - Get packet statistics for a time period +- `/api/recent_packets` - Get recent packets with all fields +- `/api/filtered_packets` - Get packets with filtering options +- `/api/packet_by_hash` - Get specific packet by hash + +### Time Series Data for Charts +- `/api/packet_type_graph_data` - Get packet type data for graphing +- `/api/metrics_graph_data` - Get metrics data for graphing +- `/api/packet_type_stats` - Get packet type distribution +- `/api/rrd_data` - Get raw RRD time series data + +## Chart.js Examples + +### 1. Packet Type Distribution (Pie Chart) + +```html + + + + + + + + + + + +``` + +### 2. Packet Types Over Time (Line Chart) + +```html + + + + + + + + + + + + +``` + +### 3. Network Metrics Dashboard (Multiple Charts) + +```html + + + + + + + + +
+
+ +
+
+ +
+
+ +
+
+ +
+
+ + + + +``` + +## Plotly.js Examples + +### 1. Interactive 3D Packet Type Surface + +```html + + + + + + +
+ + + + +``` + +### 2. Heatmap of Packet Activity + +```html + + + + + + +
+ + + + +``` + +## D3.js Example + +### Real-time Network Status + +```html + + + + + + + + + + + + +``` + +## Vue.js Framework Integration + +### Vue.js Component Example + +```vue + + + + + +``` + +### Vue.js Packet Type Dashboard + +```vue + + + + + +``` + +### Composable for API Management (Vue 3 Composition API) + +```javascript +// composables/useNetworkAPI.js +import { ref, reactive } from 'vue'; + +export function useNetworkAPI() { + const loading = ref(false); + const error = ref(null); + + const cache = reactive(new Map()); + const CACHE_TTL = 60000; // 1 minute + + async function fetchWithCache(url, forceRefresh = false) { + const cacheKey = url; + const now = Date.now(); + + // Check cache first + if (!forceRefresh && cache.has(cacheKey)) { + const cached = cache.get(cacheKey); + if (now - cached.timestamp < CACHE_TTL) { + return cached.data; + } + } + + loading.value = true; + error.value = null; + + try { + const response = await fetch(url); + const result = await response.json(); + + if (!result.success) { + throw new Error(result.error || 'API request failed'); + } + + // Cache the result + cache.set(cacheKey, { + data: result.data, + timestamp: now + }); + + return result.data; + } catch (err) { + error.value = err.message; + console.error('API Error:', err); + return null; + } finally { + loading.value = false; + } + } + + async function getPacketStats(hours = 24) { + return await fetchWithCache(`/api/packet_stats?hours=${hours}`); + } + + async function getPacketTypeStats(hours = 24) { + return await fetchWithCache(`/api/packet_type_stats?hours=${hours}`); + } + + async function getPacketTypeGraphData(hours = 24, types = 'all') { + const typesParam = types === 'all' ? '' : `&types=${types}`; + return await fetchWithCache(`/api/packet_type_graph_data?hours=${hours}${typesParam}`); + } + + async function getMetricsGraphData(hours = 24, metrics = 'all') { + const metricsParam = metrics === 'all' ? '' : `&metrics=${metrics}`; + return await fetchWithCache(`/api/metrics_graph_data?hours=${hours}${metricsParam}`); + } + + async function getRecentPackets(limit = 100) { + return await fetchWithCache(`/api/recent_packets?limit=${limit}`); + } + + async function getFilteredPackets(filters = {}) { + const params = new URLSearchParams(); + Object.entries(filters).forEach(([key, value]) => { + if (value !== null && value !== undefined) { + params.append(key, value); + } + }); + + return await fetchWithCache(`/api/filtered_packets?${params.toString()}`); + } + + function clearCache() { + cache.clear(); + } + + return { + loading, + error, + getPacketStats, + getPacketTypeStats, + getPacketTypeGraphData, + getMetricsGraphData, + getRecentPackets, + getFilteredPackets, + clearCache + }; +} +``` + +### Simple Vue.js Usage Example + +```vue + + + + + +``` +``` + +## API Usage Tips + +### 1. Error Handling +```javascript +async function fetchWithErrorHandling(url) { + try { + const response = await fetch(url); + const result = await response.json(); + + if (!result.success) { + throw new Error(result.error || 'API request failed'); + } + + return result.data; + } catch (error) { + console.error('API Error:', error); + return null; + } +} +``` + +### 2. Caching Strategy +```javascript +class APICache { + constructor(ttl = 60000) { // 1 minute TTL + this.cache = new Map(); + this.ttl = ttl; + } + + async get(key, fetcher) { + const now = Date.now(); + const cached = this.cache.get(key); + + if (cached && (now - cached.timestamp) < this.ttl) { + return cached.data; + } + + const data = await fetcher(); + this.cache.set(key, { data, timestamp: now }); + return data; + } +} + +const apiCache = new APICache(); + +// Usage +const data = await apiCache.get('metrics-24h', () => + fetch('/api/metrics_graph_data?hours=24').then(r => r.json()) +); +``` + +### 3. Real-time Updates +```javascript +function setupRealTimeUpdates(chartComponent, interval = 30000) { + const updateChart = async () => { + const data = await fetchWithErrorHandling('/api/metrics_graph_data?hours=1'); + if (data) { + chartComponent.updateData(data); + } + }; + + // Initial load + updateChart(); + + // Set up interval + const intervalId = setInterval(updateChart, interval); + + // Return cleanup function + return () => clearInterval(intervalId); +} +``` + +This documentation provides examples for creating various types of charts and visualizations using the pyMC_Repeater API endpoints. The examples cover different chart libraries and use cases, from simple statistics to real-time network monitoring. \ No newline at end of file diff --git a/repeater/http_server.py b/repeater/http_server.py deleted file mode 100644 index 02de9ca..0000000 --- a/repeater/http_server.py +++ /dev/null @@ -1,943 +0,0 @@ -import asyncio -import json -import logging -import os -import re -import threading -import time -from collections import deque -from datetime import datetime -from typing import Callable, Optional, Dict, Any - -import cherrypy -from pymc_core.protocol.utils import PAYLOAD_TYPES, ROUTE_TYPES - -from repeater import __version__ - -logger = logging.getLogger("HTTPServer") - - -# In-memory log buffer -class LogBuffer(logging.Handler): - - def __init__(self, max_lines=100): - super().__init__() - self.logs = deque(maxlen=max_lines) - self.setFormatter(logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")) - - def emit(self, record): - - try: - msg = self.format(record) - self.logs.append( - { - "message": msg, - "timestamp": datetime.fromtimestamp(record.created).isoformat(), - "level": record.levelname, - } - ) - except Exception: - self.handleError(record) - - -# Global log buffer instance -_log_buffer = LogBuffer(max_lines=100) - - -class CADCalibrationEngine: - """Real-time CAD calibration engine""" - - def __init__(self, daemon_instance=None, event_loop=None): - self.daemon_instance = daemon_instance - self.event_loop = event_loop - self.running = False - self.results = {} - self.current_test = None - self.progress = {"current": 0, "total": 0} - self.clients = set() # SSE clients - self.calibration_thread = None - - def get_test_ranges(self, spreading_factor: int): - """Get CAD test ranges""" - # Higher values = less sensitive, lower values = more sensitive - # Test from LESS sensitive to MORE sensitive to find the sweet spot - sf_ranges = { - 7: (range(22, 30, 1), range(12, 20, 1)), - 8: (range(22, 30, 1), range(12, 20, 1)), - 9: (range(24, 32, 1), range(14, 22, 1)), - 10: (range(26, 34, 1), range(16, 24, 1)), - 11: (range(28, 36, 1), range(18, 26, 1)), - 12: (range(30, 38, 1), range(20, 28, 1)), - } - return sf_ranges.get(spreading_factor, sf_ranges[8]) - - async def test_cad_config(self, radio, det_peak: int, det_min: int, samples: int = 20) -> Dict[str, Any]: - """Test CAD configuration with proper spacing and baseline measurement""" - detections = 0 - baseline_detections = 0 - - # First, get baseline with very insensitive settings (should detect nothing) - baseline_samples = 5 - for _ in range(baseline_samples): - try: - # Use very high thresholds that should detect nothing - baseline_result = await radio.perform_cad(det_peak=35, det_min=25, timeout=0.3) - if baseline_result: - baseline_detections += 1 - except Exception: - pass - await asyncio.sleep(0.1) # 100ms between baseline samples - - # Wait before actual test - await asyncio.sleep(0.5) - - # Now test the actual configuration - for i in range(samples): - try: - result = await radio.perform_cad(det_peak=det_peak, det_min=det_min, timeout=0.3) - if result: - detections += 1 - except Exception: - pass - - # Variable delay to avoid sampling artifacts - delay = 0.05 + (i % 3) * 0.05 # 50ms, 100ms, 150ms rotation - await asyncio.sleep(delay) - - # Calculate adjusted detection rate - baseline_rate = (baseline_detections / baseline_samples) * 100 - detection_rate = (detections / samples) * 100 - - # Subtract baseline noise - adjusted_rate = max(0, detection_rate - baseline_rate) - - return { - 'det_peak': det_peak, - 'det_min': det_min, - 'samples': samples, - 'detections': detections, - 'detection_rate': detection_rate, - 'baseline_rate': baseline_rate, - 'adjusted_rate': adjusted_rate, # This is the useful metric - 'sensitivity_score': self._calculate_sensitivity_score(det_peak, det_min, adjusted_rate) - } - - def _calculate_sensitivity_score(self, det_peak: int, det_min: int, adjusted_rate: float) -> float: - """Calculate a sensitivity score - higher is better balance""" - # Ideal detection rate is around 10-30% for good sensitivity without false positives - ideal_rate = 20.0 - rate_penalty = abs(adjusted_rate - ideal_rate) / ideal_rate - - # Prefer moderate sensitivity settings (not too extreme) - sensitivity_penalty = (abs(det_peak - 25) + abs(det_min - 15)) / 20.0 - - # Lower penalty = higher score - score = max(0, 100 - (rate_penalty * 50) - (sensitivity_penalty * 20)) - return score - - def broadcast_to_clients(self, data): - """Send data to all connected SSE clients""" - # Store the message for clients to pick up - self.last_message = data - # Also store in a queue for clients to consume - if not hasattr(self, 'message_queue'): - self.message_queue = [] - self.message_queue.append(data) - - def calibration_worker(self, samples: int, delay_ms: int): - """Worker thread for calibration process""" - try: - # Get radio from daemon instance - if not self.daemon_instance: - self.broadcast_to_clients({"type": "error", "message": "No daemon instance available"}) - return - - radio = getattr(self.daemon_instance, 'radio', None) - if not radio: - self.broadcast_to_clients({"type": "error", "message": "Radio instance not available"}) - return - if not hasattr(radio, 'perform_cad'): - self.broadcast_to_clients({"type": "error", "message": "Radio does not support CAD"}) - return - - # Get spreading factor from daemon instance - config = getattr(self.daemon_instance, 'config', {}) - radio_config = config.get("radio", {}) - sf = radio_config.get("spreading_factor", 8) - - # Get test ranges - peak_range, min_range = self.get_test_ranges(sf) - - total_tests = len(peak_range) * len(min_range) - self.progress = {"current": 0, "total": total_tests} - - self.broadcast_to_clients({ - "type": "status", - "message": f"Starting calibration: SF{sf}, {total_tests} tests", - "test_ranges": { - "peak_min": min(peak_range), - "peak_max": max(peak_range), - "min_min": min(min_range), - "min_max": max(min_range), - "spreading_factor": sf, - "total_tests": total_tests - } - }) - - current = 0 - - import random - - - peak_list = list(peak_range) - min_list = list(min_range) - - # Create all test combinations - test_combinations = [] - for det_peak in peak_list: - for det_min in min_list: - test_combinations.append((det_peak, det_min)) - - # Sort by distance from center for center-out pattern - peak_center = (max(peak_list) + min(peak_list)) / 2 - min_center = (max(min_list) + min(min_list)) / 2 - - def distance_from_center(combo): - peak, min_val = combo - return ((peak - peak_center) ** 2 + (min_val - min_center) ** 2) ** 0.5 - - # Sort by distance from center - test_combinations.sort(key=distance_from_center) - - - band_size = max(1, len(test_combinations) // 8) # Create 8 bands - randomized_combinations = [] - - for i in range(0, len(test_combinations), band_size): - band = test_combinations[i:i + band_size] - random.shuffle(band) # Randomize within each band - randomized_combinations.extend(band) - - # Run calibration in event loop with center-out randomized pattern - if self.event_loop: - for det_peak, det_min in randomized_combinations: - if not self.running: - break - - current += 1 - self.progress["current"] = current - - # Update progress - self.broadcast_to_clients({ - "type": "progress", - "current": current, - "total": total_tests, - "peak": det_peak, - "min": det_min - }) - - # Run the test - future = asyncio.run_coroutine_threadsafe( - self.test_cad_config(radio, det_peak, det_min, samples), - self.event_loop - ) - - try: - result = future.result(timeout=30) # 30 second timeout per test - - # Store result - key = f"{det_peak}-{det_min}" - self.results[key] = result - - # Send result to clients - self.broadcast_to_clients({ - "type": "result", - **result - }) - except Exception as e: - logger.error(f"CAD test failed for peak={det_peak}, min={det_min}: {e}") - - # Delay between tests - if self.running and delay_ms > 0: - time.sleep(delay_ms / 1000.0) - - if self.running: - # Find best result based on sensitivity score (not just detection rate) - best_result = None - recommended_result = None - if self.results: - # Find result with highest sensitivity score (best balance) - best_result = max(self.results.values(), key=lambda x: x.get('sensitivity_score', 0)) - - # Also find result with ideal adjusted detection rate (10-30%) - ideal_results = [r for r in self.results.values() if 10 <= r.get('adjusted_rate', 0) <= 30] - if ideal_results: - # Among ideal results, pick the one with best sensitivity score - recommended_result = max(ideal_results, key=lambda x: x.get('sensitivity_score', 0)) - else: - recommended_result = best_result - - self.broadcast_to_clients({ - "type": "completed", - "message": "Calibration completed", - "results": { - "best": best_result, - "recommended": recommended_result, - "total_tests": len(self.results) - } if best_result else None - }) - else: - self.broadcast_to_clients({"type": "status", "message": "Calibration stopped"}) - - except Exception as e: - logger.error(f"Calibration worker error: {e}") - self.broadcast_to_clients({"type": "error", "message": str(e)}) - finally: - self.running = False - - def start_calibration(self, samples: int = 8, delay_ms: int = 100): - """Start calibration process""" - if self.running: - return False - - self.running = True - self.results.clear() - self.progress = {"current": 0, "total": 0} - self.clear_message_queue() # Clear any old messages - - # Start calibration in separate thread - self.calibration_thread = threading.Thread( - target=self.calibration_worker, - args=(samples, delay_ms) - ) - self.calibration_thread.daemon = True - self.calibration_thread.start() - - return True - - def stop_calibration(self): - """Stop calibration process""" - self.running = False - if self.calibration_thread: - self.calibration_thread.join(timeout=2) - - def clear_message_queue(self): - """Clear the message queue when starting a new calibration""" - if hasattr(self, 'message_queue'): - self.message_queue.clear() -class APIEndpoints: - - def __init__( - self, - stats_getter: Optional[Callable] = None, - send_advert_func: Optional[Callable] = None, - config: Optional[dict] = None, - event_loop=None, - daemon_instance=None, - config_path=None, - ): - - self.stats_getter = stats_getter - self.send_advert_func = send_advert_func - self.config = config or {} - self.event_loop = event_loop - self.daemon_instance = daemon_instance - self._config_path = config_path or '/etc/pymc_repeater/config.yaml' - - # Initialize CAD calibration engine - self.cad_calibration = CADCalibrationEngine(daemon_instance, event_loop) - - @cherrypy.expose - @cherrypy.tools.json_out() - def stats(self): - - try: - stats = self.stats_getter() if self.stats_getter else {} - stats["version"] = __version__ - - # Add pyMC_Core version - try: - import pymc_core - stats["core_version"] = pymc_core.__version__ - except ImportError: - stats["core_version"] = "unknown" - - return stats - except Exception as e: - logger.error(f"Error serving stats: {e}") - return {"error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - def send_advert(self): - - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - if not self.send_advert_func: - return {"success": False, "error": "Send advert function not configured"} - - try: - import asyncio - - if self.event_loop is None: - return {"success": False, "error": "Event loop not available"} - - future = asyncio.run_coroutine_threadsafe(self.send_advert_func(), self.event_loop) - result = future.result(timeout=10) # Wait up to 10 seconds - - if result: - return {"success": True, "message": "Advert sent successfully"} - else: - return {"success": False, "error": "Failed to send advert"} - except Exception as e: - logger.error(f"Error sending advert: {e}", exc_info=True) - return {"success": False, "error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - @cherrypy.tools.json_in() - def set_mode(self): - - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - try: - data = cherrypy.request.json - new_mode = data.get("mode", "forward") - - if new_mode not in ["forward", "monitor"]: - return {"success": False, "error": "Invalid mode. Must be 'forward' or 'monitor'"} - - # Update config - if "repeater" not in self.config: - self.config["repeater"] = {} - self.config["repeater"]["mode"] = new_mode - - logger.info(f"Mode changed to: {new_mode}") - return {"success": True, "mode": new_mode} - except Exception as e: - logger.error(f"Error setting mode: {e}", exc_info=True) - return {"success": False, "error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - @cherrypy.tools.json_in() - def set_duty_cycle(self): - - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - try: - data = cherrypy.request.json - enabled = data.get("enabled", True) - - # Update config - if "duty_cycle" not in self.config: - self.config["duty_cycle"] = {} - self.config["duty_cycle"]["enforcement_enabled"] = enabled - - logger.info(f"Duty cycle enforcement {'enabled' if enabled else 'disabled'}") - return {"success": True, "enabled": enabled} - except Exception as e: - logger.error(f"Error setting duty cycle: {e}", exc_info=True) - return {"success": False, "error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - def logs(self): - - try: - logs = list(_log_buffer.logs) - return { - "logs": ( - logs - if logs - else [ - { - "message": "No logs available", - "timestamp": datetime.now().isoformat(), - "level": "INFO", - } - ] - ) - } - except Exception as e: - logger.error(f"Error fetching logs: {e}") - return {"error": str(e), "logs": []} - - # CAD Calibration endpoints - @cherrypy.expose - @cherrypy.tools.json_out() - @cherrypy.tools.json_in() - def cad_calibration_start(self): - """Start CAD calibration""" - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - try: - data = cherrypy.request.json or {} - samples = data.get("samples", 8) - delay = data.get("delay", 100) - - if self.cad_calibration.start_calibration(samples, delay): - return {"success": True, "message": "Calibration started"} - else: - return {"success": False, "error": "Calibration already running"} - - except Exception as e: - logger.error(f"Error starting CAD calibration: {e}") - return {"success": False, "error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - def cad_calibration_stop(self): - """Stop CAD calibration""" - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - try: - self.cad_calibration.stop_calibration() - return {"success": True, "message": "Calibration stopped"} - except Exception as e: - logger.error(f"Error stopping CAD calibration: {e}") - return {"success": False, "error": str(e)} - - @cherrypy.expose - @cherrypy.tools.json_out() - @cherrypy.tools.json_in() - def save_cad_settings(self): - """Save CAD calibration settings to config""" - if cherrypy.request.method != "POST": - return {"success": False, "error": "Method not allowed"} - - try: - data = cherrypy.request.json or {} - peak = data.get("peak") - min_val = data.get("min_val") - detection_rate = data.get("detection_rate", 0) - - if peak is None or min_val is None: - return {"success": False, "error": "Missing peak or min_val parameters"} - - # Update the radio immediately if available - if self.daemon_instance and hasattr(self.daemon_instance, 'radio') and self.daemon_instance.radio: - if hasattr(self.daemon_instance.radio, 'set_custom_cad_thresholds'): - self.daemon_instance.radio.set_custom_cad_thresholds(peak=peak, min_val=min_val) - logger.info(f"Applied CAD settings to radio: peak={peak}, min={min_val}") - - # Update the in-memory config - if "radio" not in self.config: - self.config["radio"] = {} - if "cad" not in self.config["radio"]: - self.config["radio"]["cad"] = {} - - self.config["radio"]["cad"]["peak_threshold"] = peak - self.config["radio"]["cad"]["min_threshold"] = min_val - - # Save to config file - config_path = getattr(self, '_config_path', '/etc/pymc_repeater/config.yaml') - self._save_config_to_file(config_path) - - logger.info(f"Saved CAD settings to config: peak={peak}, min={min_val}, rate={detection_rate:.1f}%") - return { - "success": True, - "message": f"CAD settings saved: peak={peak}, min={min_val}", - "settings": {"peak": peak, "min_val": min_val, "detection_rate": detection_rate} - } - - except Exception as e: - logger.error(f"Error saving CAD settings: {e}") - return {"success": False, "error": str(e)} - - def _save_config_to_file(self, config_path): - """Save current config to YAML file""" - try: - import yaml - import os - - # Ensure directory exists - os.makedirs(os.path.dirname(config_path), exist_ok=True) - - # Write config to file - with open(config_path, 'w') as f: - yaml.dump(self.config, f, default_flow_style=False, indent=2) - - logger.info(f"Configuration saved to {config_path}") - - except Exception as e: - logger.error(f"Failed to save config to {config_path}: {e}") - raise - - @cherrypy.expose - def cad_calibration_stream(self): - """Server-Sent Events stream for real-time updates""" - cherrypy.response.headers['Content-Type'] = 'text/event-stream' - cherrypy.response.headers['Cache-Control'] = 'no-cache' - cherrypy.response.headers['Connection'] = 'keep-alive' - cherrypy.response.headers['Access-Control-Allow-Origin'] = '*' - - def generate(): - - if not hasattr(self.cad_calibration, 'message_queue'): - self.cad_calibration.message_queue = [] - - try: - - yield f"data: {json.dumps({'type': 'connected', 'message': 'Connected to CAD calibration stream'})}\n\n" - - - if self.cad_calibration.running: - - config = getattr(self.cad_calibration.daemon_instance, 'config', {}) - radio_config = config.get("radio", {}) - sf = radio_config.get("spreading_factor", 8) - - - peak_range, min_range = self.cad_calibration.get_test_ranges(sf) - total_tests = len(peak_range) * len(min_range) - - - status_message = { - "type": "status", - "message": f"Calibration in progress: SF{sf}, {total_tests} tests", - "test_ranges": { - "peak_min": min(peak_range), - "peak_max": max(peak_range), - "min_min": min(min_range), - "min_max": max(min_range), - "spreading_factor": sf, - "total_tests": total_tests - } - } - yield f"data: {json.dumps(status_message)}\n\n" - - last_message_index = len(self.cad_calibration.message_queue) - - - while True: - - current_queue_length = len(self.cad_calibration.message_queue) - if current_queue_length > last_message_index: - - for i in range(last_message_index, current_queue_length): - message = self.cad_calibration.message_queue[i] - yield f"data: {json.dumps(message)}\n\n" - last_message_index = current_queue_length - else: - - yield f"data: {json.dumps({'type': 'keepalive'})}\n\n" - - time.sleep(0.5) - - except Exception as e: - logger.error(f"SSE stream error: {e}") - finally: - pass - - return generate() - - cad_calibration_stream._cp_config = {'response.stream': True} - - - - -class StatsApp: - - def __init__( - self, - stats_getter: Optional[Callable] = None, - template_dir: Optional[str] = None, - node_name: str = "Repeater", - pub_key: str = "", - send_advert_func: Optional[Callable] = None, - config: Optional[dict] = None, - event_loop=None, - daemon_instance=None, - config_path=None, - ): - - self.stats_getter = stats_getter - self.template_dir = template_dir - self.node_name = node_name - self.pub_key = pub_key - self.dashboard_template = None - self.config = config or {} - - # Create nested API object for routing - self.api = APIEndpoints(stats_getter, send_advert_func, self.config, event_loop, daemon_instance, config_path) - - # Load template on init - if template_dir: - template_path = os.path.join(template_dir, "dashboard.html") - try: - with open(template_path, "r") as f: - self.dashboard_template = f.read() - logger.info(f"Loaded template from {template_path}") - except FileNotFoundError: - logger.error(f"Template not found: {template_path}") - - @cherrypy.expose - def index(self): - """Serve dashboard HTML.""" - return self._serve_template("dashboard.html") - - @cherrypy.expose - def neighbors(self): - """Serve neighbors page.""" - return self._serve_template("neighbors.html") - - @cherrypy.expose - def statistics(self): - """Serve statistics page.""" - return self._serve_template("statistics.html") - - @cherrypy.expose - def configuration(self): - """Serve configuration page.""" - return self._serve_template("configuration.html") - - @cherrypy.expose - def logs(self): - """Serve logs page.""" - return self._serve_template("logs.html") - - @cherrypy.expose - def help(self): - """Serve help documentation.""" - return self._serve_template("help.html") - - @cherrypy.expose - def cad_calibration(self): - """Serve CAD calibration page.""" - return self._serve_template("cad-calibration.html") - - def _serve_template(self, template_name: str): - """Serve HTML template with stats.""" - if not self.template_dir: - return "

Error

Template directory not configured

" - - if not self.dashboard_template: - return "

Error

Template not loaded

" - - try: - - template_path = os.path.join(self.template_dir, template_name) - with open(template_path, "r") as f: - template_content = f.read() - - nav_path = os.path.join(self.template_dir, "nav.html") - nav_content = "" - try: - with open(nav_path, "r") as f: - nav_content = f.read() - except FileNotFoundError: - logger.warning(f"Navigation template not found: {nav_path}") - - stats = self.stats_getter() if self.stats_getter else {} - - if "uptime_seconds" not in stats or not isinstance( - stats.get("uptime_seconds"), (int, float) - ): - stats["uptime_seconds"] = 0 - - # Calculate uptime in hours - uptime_seconds = stats.get("uptime_seconds", 0) - uptime_hours = int(uptime_seconds // 3600) if uptime_seconds else 0 - - # Determine current page for nav highlighting - page_map = { - "dashboard.html": "dashboard", - "neighbors.html": "neighbors", - "statistics.html": "statistics", - "configuration.html": "configuration", - "cad-calibration.html": "cad-calibration", - "logs.html": "logs", - "help.html": "help", - } - current_page = page_map.get(template_name, "") - - # Prepare basic substitutions - html = template_content - html = html.replace("{{ node_name }}", str(self.node_name)) - html = html.replace("{{ last_updated }}", datetime.now().strftime("%Y-%m-%d %H:%M:%S")) - html = html.replace("{{ page }}", current_page) - - # Replace navigation placeholder with actual nav content - if "" in html: - nav_substitutions = nav_content - nav_substitutions = nav_substitutions.replace( - "{{ node_name }}", str(self.node_name) - ) - nav_substitutions = nav_substitutions.replace("{{ pub_key }}", str(self.pub_key)) - nav_substitutions = nav_substitutions.replace( - "{{ last_updated }}", datetime.now().strftime("%Y-%m-%d %H:%M:%S") - ) - - # Handle active state for nav items - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'dashboard' else '' }}", - " active" if current_page == "dashboard" else "", - ) - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'neighbors' else '' }}", - " active" if current_page == "neighbors" else "", - ) - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'statistics' else '' }}", - " active" if current_page == "statistics" else "", - ) - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'configuration' else '' }}", - " active" if current_page == "configuration" else "", - ) - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'logs' else '' }}", - " active" if current_page == "logs" else "", - ) - nav_substitutions = nav_substitutions.replace( - "{{ ' active' if page == 'help' else '' }}", - " active" if current_page == "help" else "", - ) - - html = html.replace("", nav_substitutions) - - # Build packets table HTML for dashboard - if template_name == "dashboard.html": - recent_packets = stats.get("recent_packets", []) - packets_table = "" - - if recent_packets: - for pkt in recent_packets[-20:]: # Last 20 packets - time_obj = datetime.fromtimestamp(pkt.get("timestamp", 0)) - time_str = time_obj.strftime("%H:%M:%S") - pkt_type = PAYLOAD_TYPES.get( - pkt.get("type", 0), f"0x{pkt.get('type', 0): 02x}" - ) - route_type = pkt.get("route", 0) - route = ROUTE_TYPES.get(route_type, f"UNKNOWN_{route_type}") - status = "OK TX" if pkt.get("transmitted") else "WAIT" - - # Get proper CSS class for route type - route_class = route.lower().replace("_", "-") - snr_val = pkt.get("snr", 0.0) - score_val = pkt.get("score", 0) - delay_val = pkt.get("tx_delay_ms", 0) - - packets_table += ( - "" - f"{time_str}" - f'{pkt_type}' - f'{route}' - f"{pkt.get('length', 0)}" - f"{pkt.get('rssi', 0)}" - f"{snr_val: .1f}" - f'{score_val: .2f}' - f"{delay_val: .0f}" - f"{status}" - "" - ) - else: - packets_table = """ - - - No packets received yet - waiting for traffic... - - - """ - - # Add dashboard-specific substitutions - html = html.replace("{{ rx_count }}", str(stats.get("rx_count", 0))) - html = html.replace("{{ forwarded_count }}", str(stats.get("forwarded_count", 0))) - html = html.replace("{{ dropped_count }}", str(stats.get("dropped_count", 0))) - html = html.replace("{{ uptime_hours }}", str(uptime_hours)) - - # Replace tbody with actual packets - tbody_pattern = r'.*?' - tbody_replacement = f'\n{packets_table}\n' - html = re.sub( - tbody_pattern, - tbody_replacement, - html, - flags=re.DOTALL, - ) - - return html - - except Exception as e: - logger.error(f"Error rendering template {template_name}: {e}", exc_info=True) - return f"

Error

{str(e)}

" - - -class HTTPStatsServer: - - def __init__( - self, - host: str = "0.0.0.0", - port: int = 8000, - stats_getter: Optional[Callable] = None, - template_dir: Optional[str] = None, - node_name: str = "Repeater", - pub_key: str = "", - send_advert_func: Optional[Callable] = None, - config: Optional[dict] = None, - event_loop=None, - daemon_instance=None, - config_path=None, - ): - - self.host = host - self.port = port - self.app = StatsApp( - stats_getter, template_dir, node_name, pub_key, send_advert_func, config, event_loop, daemon_instance, config_path - ) - - def start(self): - - try: - # Serve static files from templates directory - static_dir = ( - self.app.template_dir if self.app.template_dir else os.path.dirname(__file__) - ) - - config = { - "/": { - "tools.sessions.on": False, - }, - "/static": { - "tools.staticdir.on": True, - "tools.staticdir.dir": static_dir, - }, - } - - cherrypy.config.update( - { - "server.socket_host": self.host, - "server.socket_port": self.port, - "engine.autoreload.on": False, - "log.screen": False, - "log.access_file": "", # Disable access log file - "log.error_file": "", # Disable error log file - } - ) - - cherrypy.tree.mount(self.app, "/", config) - - # Completely disable access logging - cherrypy.log.access_log.propagate = False - cherrypy.log.error_log.setLevel(logging.ERROR) - - cherrypy.engine.start() - server_url = "http://{}:{}".format(self.host, self.port) - logger.info(f"HTTP stats server started on {server_url}") - - except Exception as e: - logger.error(f"Failed to start HTTP server: {e}") - raise - - def stop(self): - try: - cherrypy.engine.exit() - logger.info("HTTP stats server stopped") - except Exception as e: - logger.warning(f"Error stopping HTTP server: {e}") diff --git a/repeater/main.py b/repeater/main.py index 3de325e..65f3713 100644 --- a/repeater/main.py +++ b/repeater/main.py @@ -5,7 +5,7 @@ import sys from repeater.config import get_radio_for_board, load_config from repeater.engine import RepeaterHandler -from repeater.http_server import HTTPStatsServer, _log_buffer +from pyMC_Repeater.repeater.http.http_server import HTTPStatsServer, _log_buffer from pymc_core.node.handlers.trace import TraceHandler from pymc_core.protocol.constants import MAX_PATH_SIZE, ROUTE_TYPE_DIRECT diff --git a/repeater/storage.py b/repeater/storage.py new file mode 100644 index 0000000..306f679 --- /dev/null +++ b/repeater/storage.py @@ -0,0 +1,742 @@ +import json +import logging +import sqlite3 +import time +from pathlib import Path +from typing import Optional, Dict, Any + +try: + import rrdtool + RRDTOOL_AVAILABLE = True +except ImportError: + RRDTOOL_AVAILABLE = False + +try: + import paho.mqtt.client as mqtt + MQTT_AVAILABLE = True +except ImportError: + MQTT_AVAILABLE = False + +logger = logging.getLogger("StorageCollector") + + +class StorageCollector: + + def __init__(self, config: dict): + self.config = config + self.storage_dir = Path(config.get("storage_dir", "/var/lib/pymc_repeater")) + self.storage_dir.mkdir(parents=True, exist_ok=True) + + self.sqlite_path = self.storage_dir / "repeater.db" + self.rrd_path = self.storage_dir / "metrics.rrd" + + # MQTT configuration + self.mqtt_config = config.get("mqtt", {}) + self.mqtt_client = None + + # Initialize storage systems + self._init_sqlite() + self._init_rrd() + self._init_mqtt() + + def _init_sqlite(self): + try: + with sqlite3.connect(self.sqlite_path) as conn: + # Packets table + conn.execute(""" + CREATE TABLE IF NOT EXISTS packets ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + timestamp REAL NOT NULL, + type INTEGER NOT NULL, + route INTEGER NOT NULL, + length INTEGER NOT NULL, + rssi INTEGER, + snr REAL, + score REAL, + transmitted BOOLEAN NOT NULL, + is_duplicate BOOLEAN NOT NULL, + drop_reason TEXT, + src_hash TEXT, + dst_hash TEXT, + path_hash TEXT, + header TEXT, + payload TEXT, + payload_length INTEGER, + tx_delay_ms REAL, + packet_hash TEXT, + original_path TEXT, + forwarded_path TEXT + ) + """) + + # Adverts/neighbors table + conn.execute(""" + CREATE TABLE IF NOT EXISTS adverts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + timestamp REAL NOT NULL, + pubkey TEXT NOT NULL, + node_name TEXT, + is_repeater BOOLEAN NOT NULL, + route_type INTEGER, + contact_type TEXT, + latitude REAL, + longitude REAL, + first_seen REAL NOT NULL, + last_seen REAL NOT NULL, + rssi INTEGER, + snr REAL, + advert_count INTEGER NOT NULL DEFAULT 1, + is_new_neighbor BOOLEAN NOT NULL + ) + """) + + # Create indexes for performance + conn.execute("CREATE INDEX IF NOT EXISTS idx_packets_timestamp ON packets(timestamp)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_packets_type ON packets(type)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_packets_hash ON packets(packet_hash)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_packets_transmitted ON packets(transmitted)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_adverts_timestamp ON adverts(timestamp)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_adverts_pubkey ON adverts(pubkey)") + + conn.commit() + logger.info(f"SQLite database initialized: {self.sqlite_path}") + + except Exception as e: + logger.error(f"Failed to initialize SQLite: {e}") + + def _init_rrd(self): + + if not RRDTOOL_AVAILABLE: + logger.warning("RRDTool not available - skipping RRD initialization") + return + + if self.rrd_path.exists(): + logger.info(f"RRD database exists: {self.rrd_path}") + return + + try: + # Create RRD with 1-minute resolution, keep 1 week of detailed data + # and longer periods at reduced resolution + rrdtool.create( + str(self.rrd_path), + "--step", "60", # 1-minute steps + "--start", str(int(time.time() - 60)), + + # Data sources - Basic metrics + "DS:rx_count:COUNTER:120:0:U", # Received packets + "DS:tx_count:COUNTER:120:0:U", # Transmitted packets + "DS:drop_count:COUNTER:120:0:U", # Dropped packets + "DS:avg_rssi:GAUGE:120:-200:0", # Average RSSI + "DS:avg_snr:GAUGE:120:-30:30", # Average SNR + "DS:avg_length:GAUGE:120:0:256", # Average packet length + "DS:avg_score:GAUGE:120:0:1", # Average packet score + "DS:neighbor_count:GAUGE:120:0:U", # Number of neighbors + + # Packet type counters (based on pyMC payload types) + "DS:type_0:COUNTER:120:0:U", # Request (PAYLOAD_TYPE_REQ) + "DS:type_1:COUNTER:120:0:U", # Response (PAYLOAD_TYPE_RESPONSE) + "DS:type_2:COUNTER:120:0:U", # Text Message (PAYLOAD_TYPE_TXT_MSG) + "DS:type_3:COUNTER:120:0:U", # ACK (PAYLOAD_TYPE_ACK) + "DS:type_4:COUNTER:120:0:U", # Advert (PAYLOAD_TYPE_ADVERT) + "DS:type_5:COUNTER:120:0:U", # Group Text (PAYLOAD_TYPE_GRP_TXT) + "DS:type_6:COUNTER:120:0:U", # Group Data (PAYLOAD_TYPE_GRP_DATA) + "DS:type_7:COUNTER:120:0:U", # Anonymous Request (PAYLOAD_TYPE_ANON_REQ) + "DS:type_8:COUNTER:120:0:U", # Path (PAYLOAD_TYPE_PATH) + "DS:type_9:COUNTER:120:0:U", # Trace (PAYLOAD_TYPE_TRACE) + "DS:type_10:COUNTER:120:0:U", # Reserved for future use + "DS:type_11:COUNTER:120:0:U", # Reserved for future use + "DS:type_12:COUNTER:120:0:U", # Reserved for future use + "DS:type_13:COUNTER:120:0:U", # Reserved for future use + "DS:type_14:COUNTER:120:0:U", # Reserved for future use + "DS:type_15:COUNTER:120:0:U", # Reserved for future use + "DS:type_other:COUNTER:120:0:U", # Other packet types (>15) + + # Round Robin Archives (resolution:keep_time) + "RRA:AVERAGE:0.5:1:10080", # 1min for 1 week + "RRA:AVERAGE:0.5:5:8640", # 5min for 1 month + "RRA:AVERAGE:0.5:60:8760", # 1hour for 1 year + "RRA:MAX:0.5:1:10080", # 1min max values for 1 week + "RRA:MIN:0.5:1:10080" # 1min min values for 1 week + ) + logger.info(f"RRD database created: {self.rrd_path}") + + except Exception as e: + logger.error(f"Failed to create RRD database: {e}") + + def _init_mqtt(self): + + if not MQTT_AVAILABLE or not self.mqtt_config.get("enabled", False): + logger.info("MQTT disabled or not available") + return + + try: + self.mqtt_client = mqtt.Client() + + # Configure authentication if provided + username = self.mqtt_config.get("username") + password = self.mqtt_config.get("password") + if username: + self.mqtt_client.username_pw_set(username, password) + + # Connect to broker + broker = self.mqtt_config.get("broker", "localhost") + port = self.mqtt_config.get("port", 1883) + + self.mqtt_client.connect(broker, port, 60) + self.mqtt_client.loop_start() + + logger.info(f"MQTT client connected to {broker}:{port}") + + except Exception as e: + logger.error(f"Failed to initialize MQTT: {e}") + self.mqtt_client = None + + def record_packet(self, packet_record: dict): + + self._store_packet_sqlite(packet_record) + self._update_rrd_metrics(packet_record, record_type="packet") + self._publish_mqtt(packet_record, "packet") + + def record_advert(self, advert_record: dict): + + self._store_advert_sqlite(advert_record) + self._update_rrd_metrics(advert_record, record_type="advert") + self._publish_mqtt(advert_record, "advert") + + def _store_packet_sqlite(self, record: dict): + + try: + with sqlite3.connect(self.sqlite_path) as conn: + conn.execute(""" + INSERT INTO packets ( + timestamp, type, route, length, rssi, snr, score, + transmitted, is_duplicate, drop_reason, src_hash, dst_hash, path_hash, + header, payload, payload_length, tx_delay_ms, packet_hash, + original_path, forwarded_path + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, ( + record.get("timestamp", time.time()), + record.get("type", 0), + record.get("route", 0), + record.get("length", 0), + record.get("rssi"), + record.get("snr"), + record.get("score"), + record.get("transmitted", False), + record.get("is_duplicate", False), + record.get("drop_reason"), + record.get("src_hash"), + record.get("dst_hash"), + record.get("path_hash"), + record.get("header"), + record.get("payload"), + record.get("payload_length"), + record.get("tx_delay_ms"), + record.get("packet_hash"), + record.get("original_path"), + record.get("forwarded_path") + )) + + except Exception as e: + logger.error(f"Failed to store packet in SQLite: {e}") + + def _store_advert_sqlite(self, record: dict): + try: + with sqlite3.connect(self.sqlite_path) as conn: + # Check if this pubkey already exists + existing = conn.execute( + "SELECT pubkey, first_seen, advert_count FROM adverts WHERE pubkey = ? ORDER BY last_seen DESC LIMIT 1", + (record.get("pubkey", ""),) + ).fetchone() + + current_time = record.get("timestamp", time.time()) + + if existing: + # Update existing neighbor + conn.execute(""" + UPDATE adverts + SET timestamp = ?, node_name = ?, is_repeater = ?, route_type = ?, + contact_type = ?, latitude = ?, longitude = ?, last_seen = ?, + rssi = ?, snr = ?, advert_count = advert_count + 1, is_new_neighbor = 0 + WHERE pubkey = ? + """, ( + current_time, + record.get("node_name"), + record.get("is_repeater", False), + record.get("route_type"), + record.get("contact_type"), + record.get("latitude"), + record.get("longitude"), + current_time, + record.get("rssi"), + record.get("snr"), + record.get("pubkey", "") + )) + else: + # Insert new neighbor + conn.execute(""" + INSERT INTO adverts ( + timestamp, pubkey, node_name, is_repeater, route_type, contact_type, + latitude, longitude, first_seen, last_seen, rssi, snr, advert_count, is_new_neighbor + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, ( + current_time, + record.get("pubkey", ""), + record.get("node_name"), + record.get("is_repeater", False), + record.get("route_type"), + record.get("contact_type"), + record.get("latitude"), + record.get("longitude"), + current_time, # first_seen + current_time, # last_seen + record.get("rssi"), + record.get("snr"), + 1, # advert_count + True # is_new_neighbor + )) + + except Exception as e: + logger.error(f"Failed to store advert in SQLite: {e}") + + def _update_rrd_metrics(self, record: dict, record_type: str): + if not RRDTOOL_AVAILABLE or not self.rrd_path.exists(): + return + + try: + # Get current timestamp + timestamp = int(record.get("timestamp", time.time())) + + # Get current values from RRD (for counters we need to increment) + try: + info = rrdtool.info(str(self.rrd_path)) + last_update = int(info.get("last_update", timestamp - 60)) + + # Skip if trying to update with old data + if timestamp <= last_update: + return + + except Exception: + # If we can't read info, proceed with update + pass + + # Prepare update values based on record type + if record_type == "packet": + # Get packet type for counter tracking + packet_type = record.get("type", 0) + + # For packets, we update counters and gauges + rx_inc = 1 + tx_inc = 1 if record.get("transmitted", False) else 0 + drop_inc = 0 if record.get("transmitted", False) else 1 + + # Initialize packet type counters (all start with 0) + type_counters = ["0"] * 17 # type_0 through type_15 plus type_other + + # Increment the appropriate packet type counter + if 0 <= packet_type <= 15: + type_counters[packet_type] = "1" + else: + type_counters[16] = "1" # type_other for packet types > 15 + + # Build the values string: basic metrics + packet type counters + basic_values = f"{timestamp}:{rx_inc}:{tx_inc}:{drop_inc}:" \ + f"{record.get('rssi', 'U')}:{record.get('snr', 'U')}:" \ + f"{record.get('length', 'U')}:{record.get('score', 'U')}:U" + + type_values = ":".join(type_counters) + values = f"{basic_values}:{type_values}" + + elif record_type == "advert": + # For adverts, we mainly update gauges, packet type counters stay at 0 + type_counters = ["0"] * 17 # All packet type counters set to 0 + type_values = ":".join(type_counters) + + basic_values = f"{timestamp}:0:0:0:" \ + f"{record.get('rssi', 'U')}:{record.get('snr', 'U')}:" \ + f"U:U:1" + + values = f"{basic_values}:{type_values}" + else: + return + + rrdtool.update(str(self.rrd_path), values) + + except Exception as e: + logger.error(f"Failed to update RRD metrics: {e}") + + def _publish_mqtt(self, record: dict, record_type: str): + """Publish record to MQTT broker.""" + if not self.mqtt_client: + return + + try: + base_topic = self.mqtt_config.get("base_topic", "meshcore/repeater") + node_name = self.config.get("repeater", {}).get("node_name", "unknown") + + topic = f"{base_topic}/{node_name}/{record_type}" + + # Create clean payload (remove non-serializable items) + payload = {k: v for k, v in record.items() if v is not None} + + # Convert to JSON + message = json.dumps(payload, default=str) + + # Publish + self.mqtt_client.publish(topic, message, qos=0, retain=False) + + except Exception as e: + logger.error(f"Failed to publish to MQTT: {e}") + + def get_packet_stats(self, hours: int = 24) -> dict: + try: + cutoff = time.time() - (hours * 3600) + + with sqlite3.connect(self.sqlite_path) as conn: + conn.row_factory = sqlite3.Row + + # Basic counts + stats = conn.execute(""" + SELECT + COUNT(*) as total_packets, + SUM(transmitted) as transmitted_packets, + SUM(CASE WHEN transmitted = 0 THEN 1 ELSE 0 END) as dropped_packets, + AVG(rssi) as avg_rssi, + AVG(snr) as avg_snr, + AVG(score) as avg_score, + AVG(payload_length) as avg_payload_length, + AVG(tx_delay_ms) as avg_tx_delay + FROM packets + WHERE timestamp > ? + """, (cutoff,)).fetchone() + + # Packet types + types = conn.execute(""" + SELECT type, COUNT(*) as count + FROM packets + WHERE timestamp > ? + GROUP BY type + ORDER BY count DESC + """, (cutoff,)).fetchall() + + # Drop reasons + drop_reasons = conn.execute(""" + SELECT drop_reason, COUNT(*) as count + FROM packets + WHERE timestamp > ? AND transmitted = 0 AND drop_reason IS NOT NULL + GROUP BY drop_reason + ORDER BY count DESC + """, (cutoff,)).fetchall() + + return { + "total_packets": stats["total_packets"], + "transmitted_packets": stats["transmitted_packets"], + "dropped_packets": stats["dropped_packets"], + "avg_rssi": round(stats["avg_rssi"] or 0, 1), + "avg_snr": round(stats["avg_snr"] or 0, 1), + "avg_score": round(stats["avg_score"] or 0, 3), + "avg_payload_length": round(stats["avg_payload_length"] or 0, 1), + "avg_tx_delay": round(stats["avg_tx_delay"] or 0, 1), + "packet_types": [{"type": row["type"], "count": row["count"]} for row in types], + "drop_reasons": [{"reason": row["drop_reason"], "count": row["count"]} for row in drop_reasons] + } + + except Exception as e: + logger.error(f"Failed to get packet stats: {e}") + return {} + + def get_recent_packets(self, limit: int = 100) -> list: + """Get recent packets with all fields for debugging/analysis.""" + try: + with sqlite3.connect(self.sqlite_path) as conn: + conn.row_factory = sqlite3.Row + + packets = conn.execute(""" + SELECT + timestamp, type, route, length, rssi, snr, score, + transmitted, is_duplicate, drop_reason, src_hash, dst_hash, path_hash, + header, payload, payload_length, tx_delay_ms, packet_hash, + original_path, forwarded_path + FROM packets + ORDER BY timestamp DESC + LIMIT ? + """, (limit,)).fetchall() + + return [dict(row) for row in packets] + + except Exception as e: + logger.error(f"Failed to get recent packets: {e}") + return [] + + def get_filtered_packets(self, + packet_type: Optional[int] = None, + route: Optional[int] = None, + start_timestamp: Optional[float] = None, + end_timestamp: Optional[float] = None, + limit: int = 1000) -> list: + """Get packets filtered by type, route, and timestamp range.""" + try: + with sqlite3.connect(self.sqlite_path) as conn: + conn.row_factory = sqlite3.Row + + # Build dynamic query based on filters + where_clauses = [] + params = [] + + if packet_type is not None: + where_clauses.append("type = ?") + params.append(packet_type) + + if route is not None: + where_clauses.append("route = ?") + params.append(route) + + if start_timestamp is not None: + where_clauses.append("timestamp >= ?") + params.append(start_timestamp) + + if end_timestamp is not None: + where_clauses.append("timestamp <= ?") + params.append(end_timestamp) + + # Build the complete query + base_query = """ + SELECT + timestamp, type, route, length, rssi, snr, score, + transmitted, is_duplicate, drop_reason, src_hash, dst_hash, path_hash, + header, payload, payload_length, tx_delay_ms, packet_hash, + original_path, forwarded_path + FROM packets + """ + + if where_clauses: + query = f"{base_query} WHERE {' AND '.join(where_clauses)}" + else: + query = base_query + + query += " ORDER BY timestamp DESC LIMIT ?" + params.append(limit) + + packets = conn.execute(query, params).fetchall() + + return [dict(row) for row in packets] + + except Exception as e: + logger.error(f"Failed to get filtered packets: {e}") + return [] + + def get_packet_by_hash(self, packet_hash: str) -> Optional[dict]: + """Get a specific packet by its hash.""" + try: + with sqlite3.connect(self.sqlite_path) as conn: + conn.row_factory = sqlite3.Row + + packet = conn.execute(""" + SELECT + timestamp, type, route, length, rssi, snr, score, + transmitted, is_duplicate, drop_reason, src_hash, dst_hash, path_hash, + header, payload, payload_length, tx_delay_ms, packet_hash, + original_path, forwarded_path + FROM packets + WHERE packet_hash = ? + """, (packet_hash,)).fetchone() + + return dict(packet) if packet else None + + except Exception as e: + logger.error(f"Failed to get packet by hash: {e}") + return None + + def get_rrd_data(self, start_time: Optional[int] = None, end_time: Optional[int] = None, + resolution: str = "average") -> Optional[dict]: + """Get RRD time series data including packet type statistics.""" + if not RRDTOOL_AVAILABLE or not self.rrd_path.exists(): + return None + + try: + # Default to last 24 hours if no time specified + if end_time is None: + end_time = int(time.time()) + if start_time is None: + start_time = end_time - (24 * 3600) # 24 hours ago + + # Fetch data from RRD + fetch_result = rrdtool.fetch( + str(self.rrd_path), + resolution.upper(), + "--start", str(start_time), + "--end", str(end_time) + ) + + if not fetch_result: + return None + + (start, end, step), data_sources, data_points = fetch_result + + # Create structured response + result = { + "start_time": start, + "end_time": end, + "step": step, + "data_sources": data_sources, + "packet_types": {}, + "metrics": {} + } + + # Process data points + timestamps = [] + current_time = start + + # Initialize data arrays + for ds in data_sources: + if ds.startswith('type_'): + if 'packet_types' not in result: + result['packet_types'] = {} + result['packet_types'][ds] = [] + else: + result['metrics'][ds] = [] + + # Process each data point + for point in data_points: + timestamps.append(current_time) + + for i, value in enumerate(point): + ds_name = data_sources[i] + if ds_name.startswith('type_'): + result['packet_types'][ds_name].append(value) + else: + result['metrics'][ds_name].append(value) + + current_time += step + + result['timestamps'] = timestamps + return result + + except Exception as e: + logger.error(f"Failed to get RRD data: {e}") + return None + + def get_packet_type_stats(self, hours: int = 24) -> dict: + """Get packet type statistics for the specified time period.""" + try: + # Get RRD data for packet types + end_time = int(time.time()) + start_time = end_time - (hours * 3600) + + rrd_data = self.get_rrd_data(start_time, end_time) + if not rrd_data or 'packet_types' not in rrd_data: + return {"error": "No RRD data available"} + + # Calculate totals for each packet type + type_totals = {} + packet_type_names = { + 'type_0': 'Request (REQ)', + 'type_1': 'Response (RESPONSE)', + 'type_2': 'Text Message (TXT_MSG)', + 'type_3': 'ACK (ACK)', + 'type_4': 'Advert (ADVERT)', + 'type_5': 'Group Text (GRP_TXT)', + 'type_6': 'Group Data (GRP_DATA)', + 'type_7': 'Anonymous Request (ANON_REQ)', + 'type_8': 'Path (PATH)', + 'type_9': 'Trace (TRACE)', + 'type_10': 'Reserved Type 10', + 'type_11': 'Reserved Type 11', + 'type_12': 'Reserved Type 12', + 'type_13': 'Reserved Type 13', + 'type_14': 'Reserved Type 14', + 'type_15': 'Reserved Type 15', + 'type_other': 'Other Types (>15)' + } + + for type_key, data_points in rrd_data['packet_types'].items(): + # Calculate total (last value minus first value for counter data) + valid_points = [p for p in data_points if p is not None] + if len(valid_points) >= 2: + total = valid_points[-1] - valid_points[0] + else: + total = valid_points[0] if valid_points else 0 + + type_name = packet_type_names.get(type_key, type_key) + type_totals[type_name] = max(0, total or 0) + + return { + "hours": hours, + "packet_type_totals": type_totals, + "total_packets": sum(type_totals.values()), + "period": f"{hours} hours" + } + + except Exception as e: + logger.error(f"Failed to get packet type stats: {e}") + return {"error": str(e)} + + def get_neighbors(self) -> dict: + """Get all neighbors from the database formatted like the in-memory neighbors dict.""" + try: + with sqlite3.connect(self.sqlite_path) as conn: + conn.row_factory = sqlite3.Row + + # Get the most recent record for each pubkey + neighbors = conn.execute(""" + SELECT pubkey, node_name, is_repeater, route_type, contact_type, + latitude, longitude, first_seen, last_seen, rssi, snr, advert_count + FROM adverts a1 + WHERE last_seen = ( + SELECT MAX(last_seen) + FROM adverts a2 + WHERE a2.pubkey = a1.pubkey + ) + ORDER BY last_seen DESC + """).fetchall() + + # Convert to the same format as the in-memory neighbors dict + result = {} + for row in neighbors: + result[row["pubkey"]] = { + "node_name": row["node_name"], + "is_repeater": bool(row["is_repeater"]), + "route_type": row["route_type"], + "contact_type": row["contact_type"], + "latitude": row["latitude"], + "longitude": row["longitude"], + "first_seen": row["first_seen"], + "last_seen": row["last_seen"], + "rssi": row["rssi"], + "snr": row["snr"], + "advert_count": row["advert_count"], + } + + return result + + except Exception as e: + logger.error(f"Failed to get neighbors: {e}") + return {} + + def cleanup_old_data(self, days: int = 7): + try: + cutoff = time.time() - (days * 24 * 3600) + + with sqlite3.connect(self.sqlite_path) as conn: + # Clean old packets + result = conn.execute("DELETE FROM packets WHERE timestamp < ?", (cutoff,)) + packets_deleted = result.rowcount + + # Clean old adverts + result = conn.execute("DELETE FROM adverts WHERE timestamp < ?", (cutoff,)) + adverts_deleted = result.rowcount + + conn.commit() + + if packets_deleted > 0 or adverts_deleted > 0: + logger.info(f"Cleaned up {packets_deleted} old packets and {adverts_deleted} old adverts") + + except Exception as e: + logger.error(f"Failed to cleanup old data: {e}") + + def close(self): + """Clean shutdown of storage systems.""" + if self.mqtt_client: + self.mqtt_client.loop_stop() + self.mqtt_client.disconnect() + logger.info("MQTT client disconnected") \ No newline at end of file