diff --git a/CHANGELOG.md b/CHANGELOG.md index d83bab1..6a48dee 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,26 @@ All notable changes to SLMM (Sound Level Meter Manager) will be documented in th The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Added + +#### Wedge Detection & Automatic Recovery +- **Failure classification** - Poll failures are now probed and classified: `wedged` (control port REFUSES — device stack alive, listener gone), `offline` (timeout — power/cellular loss), or `ok` (transient). The wedge signature matches the Feb 2026 investigation: refused on 2255 while FTP/21 accepts is definitive (high confidence); refused alone is medium confidence. +- **RecoveryManager** (`app/recovery.py`) - Per-device recovery state machine: confirm wedge → trigger reset → wait for boot → reconnect → resume measurement (via existing `start_cycle` with overwrite protection) → log the full incident to device logs (category `RECOVERY`). +- **Pluggable reset backends**: + - `manual` (default) - logs/flags the wedge for a human; still auto-resumes measurement when the device returns + - `webhook` - HTTP GET/POST to a relay controller (Pi GPIO relay, Shelly, Tasmota, etc.) that power-cycles the NL-43 +- **Safety guards** - per-device `auto_recovery_enabled` master switch (default off), confirmation threshold (`WEDGE_CONFIRM_FAILURES`, default 2 consecutive failures), windowed attempt limit (default 2 per 6 h), one recovery at a time per device. +- **New endpoints**: + - `GET /api/nl43/{unit_id}/recovery/status` - connection/recovery state, wedge count, last result + - `POST /api/nl43/{unit_id}/recovery/probe` - diagnostic port probe + classification (no action taken) + - `POST /api/nl43/{unit_id}/recovery/trigger` - manually start recovery (bypasses confirmation threshold) +- **Config fields** (via `PUT /{unit_id}/config`): `auto_recovery_enabled`, `reset_backend`, `reset_webhook_url`, `reset_webhook_method`, `recovery_boot_wait_seconds`, `recovery_reconnect_timeout`, `recovery_auto_resume`, `recovery_max_attempts`, `recovery_window_minutes` +- **Status fields**: `connection_state`, `last_wedge_at`, `wedge_count`, `recovery_state`, `last_recovery_at`, `last_recovery_result` +- **Migration**: `migrate_add_recovery_fields.py` +- **Test suite**: `test_wedge_recovery.py` - fake wedge-able NL-43 + fake relay webhook; covers classification, end-to-end recovery with measurement resume, and gating (23 assertions) + ## [0.3.0] - 2026-02-17 ### Added diff --git a/app/background_poller.py b/app/background_poller.py index 64bcc60..b63bde1 100644 --- a/app/background_poller.py +++ b/app/background_poller.py @@ -17,6 +17,7 @@ from app.database import SessionLocal from app.models import NL43Config, NL43Status from app.services import NL43Client, persist_snapshot, sync_measurement_start_time_from_ftp from app.device_logger import log_device_event, cleanup_old_logs +from app.recovery import classify_failure, recovery_manager logger = logging.getLogger(__name__) @@ -231,6 +232,7 @@ class BackgroundPoller: status.consecutive_failures = 0 status.last_success = datetime.utcnow() status.last_error = None + status.connection_state = "ok" db.commit() self._logger.info(f"✓ Successfully polled {unit_id}") @@ -322,6 +324,61 @@ class BackgroundPoller: db.commit() + # Classify the failure: wedge (listener dead, host alive) vs offline + await self._classify_and_recover(cfg, status, db) + + async def _classify_and_recover(self, cfg: NL43Config, status: NL43Status, db: Session): + """Probe the device to classify a poll failure and trigger wedge recovery. + + The wedge signature (see SLM-stress-test investigation doc): control + port REFUSES connections (RST = host stack alive, listener gone) while + the device is otherwise up. A timeout means offline (power/cell loss). + """ + unit_id = cfg.unit_id + + # Skip while a recovery for this device is already running — its own + # probes manage state until it finishes. + if recovery_manager.is_active(unit_id): + return + + try: + classification = await classify_failure( + cfg.host, + cfg.tcp_port, + ftp_port=(cfg.ftp_port or 21) if cfg.ftp_enabled else None, + ) + except Exception as e: + self._logger.warning(f"Failure classification error for {unit_id}: {e}") + return + + new_state = classification["state"] + prev_state = status.connection_state or "ok" + + if new_state != prev_state: + status.connection_state = new_state + if new_state == "wedged": + status.last_wedge_at = datetime.utcnow() + status.wedge_count = (status.wedge_count or 0) + 1 + log_device_event( + unit_id, "ERROR", "RECOVERY", + f"WEDGE detected: control port {cfg.tcp_port} refused " + f"(tcp_probe={classification['tcp_probe']}, ftp_probe={classification['ftp_probe']}, " + f"confidence={classification['confidence']})", + db, + ) + self._logger.error(f"Device {unit_id} classified as WEDGED ({classification})") + elif new_state == "offline": + log_device_event( + unit_id, "WARNING", "RECOVERY", + f"Device classified OFFLINE (tcp_probe={classification['tcp_probe']}) - " + f"not a wedge, no reset will be triggered", + db, + ) + db.commit() + + if new_state == "wedged": + recovery_manager.maybe_start_recovery(cfg, status, classification, db) + def _calculate_sleep_interval(self) -> int: """ Calculate the next sleep interval based on all device poll intervals. diff --git a/app/models.py b/app/models.py index 4c86514..2871e85 100644 --- a/app/models.py +++ b/app/models.py @@ -23,6 +23,17 @@ class NL43Config(Base): poll_interval_seconds = Column(Integer, nullable=True, default=60) # Polling interval (10-3600 seconds) poll_enabled = Column(Boolean, default=True) # Enable/disable background polling for this device + # Wedge recovery configuration + auto_recovery_enabled = Column(Boolean, default=False) # Master switch for automatic wedge recovery + reset_backend = Column(String, default="manual") # Reset actuation: "manual" (notify only) or "webhook" + reset_webhook_url = Column(String, nullable=True) # URL that triggers a power cycle (relay controller) + reset_webhook_method = Column(String, default="POST") # HTTP method for webhook (GET or POST) + recovery_boot_wait_seconds = Column(Integer, default=90) # Wait after reset before reconnect attempts + recovery_reconnect_timeout = Column(Integer, default=180) # Max seconds to wait for device after boot wait + recovery_auto_resume = Column(Boolean, default=True) # Restart measurement if it was running pre-wedge + recovery_max_attempts = Column(Integer, default=2) # Max recovery attempts per window + recovery_window_minutes = Column(Integer, default=360) # Attempt-limit window (default 6 hours) + class NL43Status(Base): """ @@ -57,6 +68,14 @@ class NL43Status(Base): # FTP start time sync tracking start_time_sync_attempted = Column(Boolean, default=False) # True if FTP sync was attempted for current measurement + # Wedge detection and recovery tracking + connection_state = Column(String, default="ok") # ok | wedged | degraded | offline + last_wedge_at = Column(DateTime, nullable=True) # When the current/most recent wedge was detected + wedge_count = Column(Integer, default=0) # Lifetime count of detected wedges + recovery_state = Column(String, default="idle") # idle | resetting | waiting_boot | reconnecting | resuming | awaiting_manual_reset | failed + last_recovery_at = Column(DateTime, nullable=True) # When the last recovery attempt started + last_recovery_result = Column(Text, nullable=True) # Outcome summary of the last recovery attempt + class DeviceLog(Base): """ diff --git a/app/recovery.py b/app/recovery.py new file mode 100644 index 0000000..837d3a7 --- /dev/null +++ b/app/recovery.py @@ -0,0 +1,482 @@ +""" +Wedge detection and automatic recovery for NL-43 devices. + +The NL-43's TCP control listener (port 2255) can stop listening while the +device otherwise remains healthy (LAN stack alive, FTP alive) — the "wedge" +failure mode documented in SLM-stress-test/NL43_RX55_TCP_Wedge_Investigation_2026-02-18.md. +Only a power cycle of the NL-43 has been proven to restore the listener. + +This module provides: + +1. Failure classification — distinguishes a wedge (connection REFUSED, i.e. + the device's IP stack answered with RST so the host is alive but the + listener is gone) from the device being offline (connect timeout — power + loss, cellular outage). An FTP-port probe adds confidence when available. + +2. RecoveryManager — a per-device state machine that, once a wedge is + confirmed, triggers a reset through a pluggable backend, waits for the + device to boot, reconnects, and resumes measurement if one was running. + +Reset backends: +- "manual" — log/flag only; a human power-cycles the device. The manager + keeps watching and still auto-resumes measurement when the + device returns. +- "webhook" — HTTP request to a relay controller (Pi GPIO relay, Shelly, + Tasmota, RX55 GPIO bridge, etc.) that power-cycles the device. + +Recovery states (NL43Status.recovery_state): + idle -> resetting -> waiting_boot -> reconnecting -> resuming -> idle + \\-> awaiting_manual_reset -> reconnecting -> ... + any step may end in "failed" (recorded in last_recovery_result). +""" + +import asyncio +import json +import logging +import os +import time +import urllib.request +from datetime import datetime +from typing import Dict, List, Optional + +from sqlalchemy.orm import Session + +from app.database import SessionLocal +from app.models import NL43Config, NL43Status +from app.device_logger import log_device_event + +logger = logging.getLogger(__name__) + +# Number of consecutive wedge-classified poll failures required before +# automatic recovery is triggered (manual trigger bypasses this). +WEDGE_CONFIRM_FAILURES = int(os.getenv("WEDGE_CONFIRM_FAILURES", "2")) + +# Timeout for individual TCP probes during classification/reconnection +PROBE_TIMEOUT = float(os.getenv("RECOVERY_PROBE_TIMEOUT", "5.0")) + +# Seconds between reconnect probes after a reset +RECONNECT_PROBE_INTERVAL = 10 +# Seconds between probes while awaiting a manual reset +MANUAL_WATCH_PROBE_INTERVAL = 30 + + +# --------------------------------------------------------------------------- +# TCP probing and failure classification +# --------------------------------------------------------------------------- + +async def probe_tcp(host: str, port: int, timeout: float = PROBE_TIMEOUT) -> str: + """Probe a TCP port and report what happened. + + Returns one of: + "open" — connect succeeded (closed again immediately) + "refused" — RST received: host stack alive, no listener + "timeout" — no response: host down / network unreachable / filtered + "unreachable" — OS-level error (no route, etc.) + """ + try: + reader, writer = await asyncio.wait_for( + asyncio.open_connection(host, port), timeout=timeout + ) + writer.close() + try: + await writer.wait_closed() + except Exception: + pass + return "open" + except asyncio.TimeoutError: + return "timeout" + except ConnectionRefusedError: + return "refused" + except OSError: + return "unreachable" + + +async def classify_failure(host: str, tcp_port: int, ftp_port: Optional[int] = 21) -> dict: + """Classify a poll failure as wedged / offline / ok. + + The wedge signature (from the Feb 2026 investigation): the control port + actively REFUSES (RST) while the device is otherwise alive. A refused + connection requires a live IP stack at the far end — a powered-off device + behind the modem produces a timeout instead. The FTP probe upgrades + confidence when it accepts, but refusal of the control port alone is + already a strong wedge indicator. + + Returns dict: + state — "ok" | "wedged" | "offline" + confidence — "high" | "medium" (only meaningful for "wedged") + tcp_probe — raw probe result for the control port + ftp_probe — raw probe result for the FTP port (None if not probed) + """ + tcp_result = await probe_tcp(host, tcp_port) + + result = { + "state": "ok", + "confidence": "high", + "tcp_probe": tcp_result, + "ftp_probe": None, + } + + if tcp_result == "open": + # Transient failure — the listener is there now + return result + + if tcp_result in ("timeout", "unreachable"): + result["state"] = "offline" + return result + + # tcp_result == "refused": host alive, control listener gone + result["state"] = "wedged" + if ftp_port: + ftp_result = await probe_tcp(host, ftp_port) + result["ftp_probe"] = ftp_result + # FTP accepting while control refuses is the definitive signature + result["confidence"] = "high" if ftp_result == "open" else "medium" + else: + result["confidence"] = "medium" + + return result + + +# --------------------------------------------------------------------------- +# Reset backends +# --------------------------------------------------------------------------- + +class ResetBackend: + """Base class for reset actuation backends. + + Backends receive a plain dict of snapshotted config values (never a live + ORM row — recovery runs as a detached task after the originating DB + session has closed). + """ + + name = "base" + + async def trigger(self, params: dict) -> dict: + """Attempt to power-cycle the device. + + Returns dict: {"triggered": bool, "detail": str} + "triggered" True means a reset is expected to be in progress. + """ + raise NotImplementedError + + +class ManualNotifyBackend(ResetBackend): + """No actuation — flags the wedge for a human to power-cycle the device.""" + + name = "manual" + + async def trigger(self, params: dict) -> dict: + return { + "triggered": False, + "detail": "No reset hardware configured - manual power cycle required", + } + + +class WebhookBackend(ResetBackend): + """Calls an HTTP endpoint that power-cycles the device. + + The endpoint is expected to perform the full cycle itself (e.g. a Shelly + relay with `?turn=off&timer=10`, a Tasmota PulseTime rule, or a small + HTTP service on a Pi driving a relay GPIO). Any 2xx response counts as + triggered. + """ + + name = "webhook" + + async def trigger(self, params: dict) -> dict: + url = params.get("webhook_url") + if not url: + return {"triggered": False, "detail": "reset_backend is 'webhook' but reset_webhook_url is not set"} + + method = (params.get("webhook_method") or "POST").upper() + + def _call() -> dict: + req = urllib.request.Request(url, method=method) + req.add_header("User-Agent", "SLMM-wedge-recovery") + with urllib.request.urlopen(req, timeout=15) as resp: + body = resp.read(500).decode(errors="ignore") + return {"status": resp.status, "body": body} + + try: + resp = await asyncio.to_thread(_call) + if 200 <= resp["status"] < 300: + return {"triggered": True, "detail": f"Webhook {method} {url} -> {resp['status']}"} + return {"triggered": False, "detail": f"Webhook {method} {url} -> HTTP {resp['status']}"} + except Exception as e: + return {"triggered": False, "detail": f"Webhook {method} {url} failed: {e}"} + + +BACKENDS: Dict[str, ResetBackend] = { + "manual": ManualNotifyBackend(), + "webhook": WebhookBackend(), +} + + +def get_backend(name: Optional[str]) -> ResetBackend: + return BACKENDS.get((name or "manual").lower(), BACKENDS["manual"]) + + +# --------------------------------------------------------------------------- +# Recovery manager +# --------------------------------------------------------------------------- + +class RecoveryManager: + """Orchestrates wedge recovery, one concurrent recovery per device.""" + + def __init__(self): + self._active: Dict[str, asyncio.Task] = {} + # In-memory attempt timestamps per unit for windowed rate limiting. + # Resets on service restart; last_recovery_at in the DB provides a + # coarse cross-restart guard via the same window check. + self._attempts: Dict[str, List[float]] = {} + + # -- public API --------------------------------------------------------- + + def is_active(self, unit_id: str) -> bool: + task = self._active.get(unit_id) + return task is not None and not task.done() + + def get_status(self, unit_id: str) -> dict: + now = time.time() + attempts = self._attempts.get(unit_id, []) + return { + "recovery_in_progress": self.is_active(unit_id), + "attempts_this_session": len(attempts), + "last_attempt_seconds_ago": round(now - attempts[-1]) if attempts else None, + } + + def maybe_start_recovery(self, cfg: NL43Config, status: NL43Status, classification: dict, db: Session) -> bool: + """Start a recovery task if conditions are met. Called from the poller. + + Returns True if a recovery task was started. + """ + unit_id = cfg.unit_id + + if not cfg.auto_recovery_enabled: + logger.info(f"[RECOVERY] {unit_id}: wedge detected but auto_recovery_enabled is off") + return False + + if self.is_active(unit_id): + logger.debug(f"[RECOVERY] {unit_id}: recovery already in progress") + return False + + if status.consecutive_failures < WEDGE_CONFIRM_FAILURES: + logger.info( + f"[RECOVERY] {unit_id}: wedge suspected " + f"({status.consecutive_failures}/{WEDGE_CONFIRM_FAILURES} confirmations)" + ) + return False + + if not self._attempt_allowed(cfg, status): + logger.warning( + f"[RECOVERY] {unit_id}: attempt limit reached " + f"({cfg.recovery_max_attempts} per {cfg.recovery_window_minutes} min) - standing down" + ) + log_device_event( + unit_id, "ERROR", "RECOVERY", + f"Wedge persists but recovery attempt limit reached " + f"({cfg.recovery_max_attempts} per {cfg.recovery_window_minutes} min). Manual intervention required.", + db, + ) + return False + + return self.start_recovery(cfg, classification) + + def start_recovery(self, cfg: NL43Config, classification: Optional[dict] = None) -> bool: + """Launch the recovery task (also used by the manual-trigger endpoint).""" + unit_id = cfg.unit_id + if self.is_active(unit_id): + return False + + self._attempts.setdefault(unit_id, []).append(time.time()) + + # Snapshot config values so the task doesn't depend on a live DB row + params = { + "unit_id": unit_id, + "host": cfg.host, + "tcp_port": cfg.tcp_port, + "ftp_port": cfg.ftp_port or 21, + "ftp_enabled": bool(cfg.ftp_enabled), + "ftp_username": cfg.ftp_username, + "ftp_password": cfg.ftp_password, + "backend": (cfg.reset_backend or "manual").lower(), + "webhook_url": cfg.reset_webhook_url, + "webhook_method": cfg.reset_webhook_method, + "boot_wait": cfg.recovery_boot_wait_seconds or 90, + "reconnect_timeout": cfg.recovery_reconnect_timeout or 180, + "auto_resume": bool(cfg.recovery_auto_resume), + "window_minutes": cfg.recovery_window_minutes or 360, + "classification": classification or {}, + } + + task = asyncio.create_task(self._run_recovery(params)) + self._active[unit_id] = task + logger.info(f"[RECOVERY] {unit_id}: recovery task started (backend={params['backend']})") + return True + + # -- internals ---------------------------------------------------------- + + def _attempt_allowed(self, cfg: NL43Config, status: NL43Status) -> bool: + """Windowed attempt limit: max N attempts per window.""" + window_seconds = (cfg.recovery_window_minutes or 360) * 60 + max_attempts = cfg.recovery_max_attempts or 2 + now = time.time() + + recent = [t for t in self._attempts.get(cfg.unit_id, []) if now - t < window_seconds] + self._attempts[cfg.unit_id] = recent + if len(recent) >= max_attempts: + return False + + # Cross-restart guard: if the in-memory history is empty (service + # restarted) but the DB shows a very recent attempt, respect a + # cooldown of one window-fraction to avoid rapid-fire cycling. + if not recent and status.last_recovery_at is not None: + elapsed = (datetime.utcnow() - status.last_recovery_at).total_seconds() + if elapsed < window_seconds / max_attempts / 2: + return False + + return True + + def _set_state(self, unit_id: str, recovery_state: str, result: Optional[str] = None, + connection_state: Optional[str] = None, mark_recovered: bool = False): + """Persist recovery progress on the status row (own short-lived session).""" + db = SessionLocal() + try: + status = db.query(NL43Status).filter_by(unit_id=unit_id).first() + if not status: + status = NL43Status(unit_id=unit_id) + db.add(status) + status.recovery_state = recovery_state + if result is not None: + status.last_recovery_result = result[:1000] + if connection_state is not None: + status.connection_state = connection_state + if mark_recovered: + status.is_reachable = True + status.consecutive_failures = 0 + db.commit() + except Exception as e: + logger.warning(f"[RECOVERY] {unit_id}: failed to persist state '{recovery_state}': {e}") + finally: + db.close() + + async def _run_recovery(self, p: dict): + """The recovery state machine. Runs as a detached task.""" + unit_id = p["unit_id"] + backend = get_backend(p["backend"]) + started = time.time() + + def _log(level: str, message: str): + log_device_event(unit_id, level, "RECOVERY", message) + + try: + # Record attempt start + db = SessionLocal() + try: + status = db.query(NL43Status).filter_by(unit_id=unit_id).first() + want_resume = bool(status and status.measurement_state == "Start") + if status: + status.last_recovery_at = datetime.utcnow() + status.recovery_state = "resetting" + db.commit() + finally: + db.close() + + cls = p.get("classification") or {} + _log("ERROR", + f"Wedge recovery started (backend={backend.name}, " + f"tcp_probe={cls.get('tcp_probe', '?')}, ftp_probe={cls.get('ftp_probe', '?')}, " + f"was_measuring={want_resume})") + + # --- Step 1: trigger reset --------------------------------------- + trigger_result = await backend.trigger(p) + _log("INFO", f"Reset trigger: {trigger_result['detail']}") + + if trigger_result["triggered"]: + # --- Step 2: wait for the device to boot --------------------- + self._set_state(unit_id, "waiting_boot") + _log("INFO", f"Reset triggered - waiting {p['boot_wait']}s for device boot") + await asyncio.sleep(p["boot_wait"]) + reconnect_deadline = time.time() + p["reconnect_timeout"] + probe_interval = RECONNECT_PROBE_INTERVAL + else: + if backend.name == "manual": + # Watch for a human to power-cycle the device + self._set_state(unit_id, "awaiting_manual_reset") + _log("ERROR", + f"WEDGE DETECTED on {p['host']}:{p['tcp_port']} - " + f"manual power cycle required. Watching for device return " + f"(up to {p['window_minutes']} min).") + reconnect_deadline = time.time() + p["window_minutes"] * 60 + probe_interval = MANUAL_WATCH_PROBE_INTERVAL + else: + raise RuntimeError(f"Reset trigger failed: {trigger_result['detail']}") + + # --- Step 3: wait for the control port to come back -------------- + self._set_state(unit_id, "reconnecting") + port_open = False + while time.time() < reconnect_deadline: + probe = await probe_tcp(p["host"], p["tcp_port"]) + if probe == "open": + port_open = True + break + await asyncio.sleep(probe_interval) + + if not port_open: + raise RuntimeError( + f"Device did not return within " + f"{round(reconnect_deadline - started)}s of recovery start" + ) + + _log("INFO", f"Control port {p['tcp_port']} is accepting connections again") + + # --- Step 4: flush stale pooled connection, verify protocol ------ + from app.services import NL43Client, _connection_pool + + device_key = f"{p['host']}:{p['tcp_port']}" + await _connection_pool.discard(device_key) + + client = NL43Client( + p["host"], p["tcp_port"], timeout=10.0, + ftp_username=p["ftp_username"], ftp_password=p["ftp_password"], + ftp_port=p["ftp_port"], + ) + state = await client.get_measurement_state() + _log("INFO", f"Device responding to commands - measurement state: {state}") + + # --- Step 5: resume measurement if needed ------------------------- + resumed = False + if want_resume and p["auto_resume"] and state != "Start": + self._set_state(unit_id, "resuming") + _log("INFO", "Measurement was running before the wedge - executing start cycle") + cycle = await client.start_cycle(sync_clock=True) + resumed = True + _log("INFO", + f"Measurement restarted (index {cycle.get('old_index')} -> {cycle.get('new_index')})") + elif want_resume and state == "Start": + _log("INFO", "Device is already measuring after reset - no resume needed") + resumed = True + + # --- Done --------------------------------------------------------- + elapsed = round(time.time() - started) + summary = ( + f"Recovered in {elapsed}s via {backend.name}" + + (", measurement resumed" if resumed else + (", measurement NOT resumed" if want_resume else "")) + ) + self._set_state(unit_id, "idle", result=summary, connection_state="ok", mark_recovered=True) + _log("INFO", f"Wedge recovery SUCCEEDED: {summary}") + + except Exception as e: + elapsed = round(time.time() - started) + summary = f"Recovery FAILED after {elapsed}s: {e}" + self._set_state(unit_id, "failed", result=summary) + _log("ERROR", summary) + logger.error(f"[RECOVERY] {unit_id}: {summary}") + + finally: + self._active.pop(unit_id, None) + + +# Global singleton +recovery_manager = RecoveryManager() diff --git a/app/routers.py b/app/routers.py index a21c928..b371443 100644 --- a/app/routers.py +++ b/app/routers.py @@ -13,6 +13,7 @@ import asyncio from app.database import get_db from app.models import NL43Config, NL43Status from app.services import NL43Client, persist_snapshot +from app.recovery import classify_failure, recovery_manager logger = logging.getLogger(__name__) @@ -53,6 +54,38 @@ class ConfigPayload(BaseModel): poll_enabled: bool | None = None poll_interval_seconds: int | None = None + # Wedge recovery configuration + auto_recovery_enabled: bool | None = None + reset_backend: str | None = None + reset_webhook_url: str | None = None + reset_webhook_method: str | None = None + recovery_boot_wait_seconds: int | None = Field(None, ge=10, le=600) + recovery_reconnect_timeout: int | None = Field(None, ge=30, le=3600) + recovery_auto_resume: bool | None = None + recovery_max_attempts: int | None = Field(None, ge=1, le=10) + recovery_window_minutes: int | None = Field(None, ge=30, le=1440) + + @field_validator("reset_backend") + @classmethod + def validate_reset_backend(cls, v): + if v is not None and v.lower() not in ("manual", "webhook"): + raise ValueError("reset_backend must be 'manual' or 'webhook'") + return v.lower() if v else v + + @field_validator("reset_webhook_method") + @classmethod + def validate_reset_webhook_method(cls, v): + if v is not None and v.upper() not in ("GET", "POST"): + raise ValueError("reset_webhook_method must be 'GET' or 'POST'") + return v.upper() if v else v + + @field_validator("reset_webhook_url") + @classmethod + def validate_reset_webhook_url(cls, v): + if v is not None and v != "" and not v.lower().startswith(("http://", "https://")): + raise ValueError("reset_webhook_url must start with http:// or https://") + return v + @field_validator("host") @classmethod def validate_host(cls, v): @@ -348,6 +381,15 @@ def get_config(unit_id: str, db: Session = Depends(get_db)): "ftp_username": cfg.ftp_username, "ftp_password": cfg.ftp_password, "web_enabled": cfg.web_enabled, + "auto_recovery_enabled": cfg.auto_recovery_enabled, + "reset_backend": cfg.reset_backend, + "reset_webhook_url": cfg.reset_webhook_url, + "reset_webhook_method": cfg.reset_webhook_method, + "recovery_boot_wait_seconds": cfg.recovery_boot_wait_seconds, + "recovery_reconnect_timeout": cfg.recovery_reconnect_timeout, + "recovery_auto_resume": cfg.recovery_auto_resume, + "recovery_max_attempts": cfg.recovery_max_attempts, + "recovery_window_minutes": cfg.recovery_window_minutes, }, } @@ -407,6 +449,24 @@ async def upsert_config(unit_id: str, payload: ConfigPayload, db: Session = Depe cfg.poll_enabled = payload.poll_enabled if payload.poll_interval_seconds is not None: cfg.poll_interval_seconds = payload.poll_interval_seconds + if payload.auto_recovery_enabled is not None: + cfg.auto_recovery_enabled = payload.auto_recovery_enabled + if payload.reset_backend is not None: + cfg.reset_backend = payload.reset_backend + if payload.reset_webhook_url is not None: + cfg.reset_webhook_url = payload.reset_webhook_url or None + if payload.reset_webhook_method is not None: + cfg.reset_webhook_method = payload.reset_webhook_method + if payload.recovery_boot_wait_seconds is not None: + cfg.recovery_boot_wait_seconds = payload.recovery_boot_wait_seconds + if payload.recovery_reconnect_timeout is not None: + cfg.recovery_reconnect_timeout = payload.recovery_reconnect_timeout + if payload.recovery_auto_resume is not None: + cfg.recovery_auto_resume = payload.recovery_auto_resume + if payload.recovery_max_attempts is not None: + cfg.recovery_max_attempts = payload.recovery_max_attempts + if payload.recovery_window_minutes is not None: + cfg.recovery_window_minutes = payload.recovery_window_minutes db.commit() db.refresh(cfg) @@ -803,6 +863,89 @@ async def get_measurement_state(unit_id: str, db: Session = Depends(get_db)): raise HTTPException(status_code=502, detail=str(e)) +# ============================================================================ +# WEDGE RECOVERY ENDPOINTS +# ============================================================================ + +@router.get("/{unit_id}/recovery/status") +def get_recovery_status(unit_id: str, db: Session = Depends(get_db)): + """Get wedge detection and recovery status for a device.""" + cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first() + if not cfg: + raise HTTPException(status_code=404, detail="NL43 config not found") + + status = db.query(NL43Status).filter_by(unit_id=unit_id).first() + manager_status = recovery_manager.get_status(unit_id) + + return { + "status": "ok", + "data": { + "unit_id": unit_id, + "connection_state": status.connection_state if status else "unknown", + "recovery_state": status.recovery_state if status else "idle", + "last_wedge_at": status.last_wedge_at.isoformat() + "Z" if status and status.last_wedge_at else None, + "wedge_count": status.wedge_count if status else 0, + "last_recovery_at": status.last_recovery_at.isoformat() + "Z" if status and status.last_recovery_at else None, + "last_recovery_result": status.last_recovery_result if status else None, + "auto_recovery_enabled": cfg.auto_recovery_enabled, + "reset_backend": cfg.reset_backend, + **manager_status, + }, + } + + +@router.post("/{unit_id}/recovery/probe") +async def probe_device(unit_id: str, db: Session = Depends(get_db)): + """Probe the device's control and FTP ports and classify its state. + + Does not trigger recovery — diagnostic only. Useful for confirming a + wedge before manually triggering a reset. + """ + cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first() + if not cfg: + raise HTTPException(status_code=404, detail="NL43 config not found") + + classification = await classify_failure( + cfg.host, + cfg.tcp_port, + ftp_port=(cfg.ftp_port or 21) if cfg.ftp_enabled else None, + ) + return {"status": "ok", "unit_id": unit_id, "data": classification} + + +@router.post("/{unit_id}/recovery/trigger") +async def trigger_recovery(unit_id: str, db: Session = Depends(get_db)): + """Manually trigger the recovery flow for a device. + + Bypasses the consecutive-failure confirmation threshold but still + respects the one-recovery-at-a-time guard. Use after confirming a wedge + via /recovery/probe, or to test the reset hardware end-to-end. + """ + cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first() + if not cfg: + raise HTTPException(status_code=404, detail="NL43 config not found") + + if recovery_manager.is_active(unit_id): + raise HTTPException(status_code=409, detail="Recovery already in progress for this device") + + classification = await classify_failure( + cfg.host, + cfg.tcp_port, + ftp_port=(cfg.ftp_port or 21) if cfg.ftp_enabled else None, + ) + + started = recovery_manager.start_recovery(cfg, classification) + if not started: + raise HTTPException(status_code=409, detail="Could not start recovery (already in progress)") + + return { + "status": "ok", + "unit_id": unit_id, + "message": f"Recovery started (backend={cfg.reset_backend or 'manual'})", + "classification": classification, + } + + @router.post("/{unit_id}/sleep") async def sleep_device(unit_id: str, db: Session = Depends(get_db)): """Put the device into sleep mode for battery conservation.""" diff --git a/migrate_add_recovery_fields.py b/migrate_add_recovery_fields.py new file mode 100644 index 0000000..2bd624e --- /dev/null +++ b/migrate_add_recovery_fields.py @@ -0,0 +1,107 @@ +#!/usr/bin/env python3 +""" +Migration script to add wedge-recovery fields to nl43_config and nl43_status tables. + +Adds to nl43_config: +- auto_recovery_enabled (BOOLEAN, default 0/False) +- reset_backend (TEXT, default 'manual') +- reset_webhook_url (TEXT, nullable) +- reset_webhook_method (TEXT, default 'POST') +- recovery_boot_wait_seconds (INTEGER, default 90) +- recovery_reconnect_timeout (INTEGER, default 180) +- recovery_auto_resume (BOOLEAN, default 1/True) +- recovery_max_attempts (INTEGER, default 2) +- recovery_window_minutes (INTEGER, default 360) + +Adds to nl43_status: +- connection_state (TEXT, default 'ok') +- last_wedge_at (DATETIME, nullable) +- wedge_count (INTEGER, default 0) +- recovery_state (TEXT, default 'idle') +- last_recovery_at (DATETIME, nullable) +- last_recovery_result (TEXT, nullable) + +Usage: + python migrate_add_recovery_fields.py +""" + +import sqlite3 +import sys +from pathlib import Path + +CONFIG_COLUMNS = [ + ("auto_recovery_enabled", "BOOLEAN DEFAULT 0"), + ("reset_backend", "TEXT DEFAULT 'manual'"), + ("reset_webhook_url", "TEXT"), + ("reset_webhook_method", "TEXT DEFAULT 'POST'"), + ("recovery_boot_wait_seconds", "INTEGER DEFAULT 90"), + ("recovery_reconnect_timeout", "INTEGER DEFAULT 180"), + ("recovery_auto_resume", "BOOLEAN DEFAULT 1"), + ("recovery_max_attempts", "INTEGER DEFAULT 2"), + ("recovery_window_minutes", "INTEGER DEFAULT 360"), +] + +STATUS_COLUMNS = [ + ("connection_state", "TEXT DEFAULT 'ok'"), + ("last_wedge_at", "DATETIME"), + ("wedge_count", "INTEGER DEFAULT 0"), + ("recovery_state", "TEXT DEFAULT 'idle'"), + ("last_recovery_at", "DATETIME"), + ("last_recovery_result", "TEXT"), +] + + +def migrate(): + db_path = Path("data/slmm.db") + + if not db_path.exists(): + print(f"❌ Database not found at {db_path}") + print(" Run this script from the slmm directory") + return False + + try: + conn = sqlite3.connect(db_path) + cursor = conn.cursor() + + cursor.execute("PRAGMA table_info(nl43_config)") + config_columns = [row[1] for row in cursor.fetchall()] + + cursor.execute("PRAGMA table_info(nl43_status)") + status_columns = [row[1] for row in cursor.fetchall()] + + changes_made = False + + for name, ddl in CONFIG_COLUMNS: + if name not in config_columns: + print(f"Adding {name} to nl43_config...") + cursor.execute(f"ALTER TABLE nl43_config ADD COLUMN {name} {ddl}") + changes_made = True + else: + print(f"✓ {name} already exists in nl43_config") + + for name, ddl in STATUS_COLUMNS: + if name not in status_columns: + print(f"Adding {name} to nl43_status...") + cursor.execute(f"ALTER TABLE nl43_status ADD COLUMN {name} {ddl}") + changes_made = True + else: + print(f"✓ {name} already exists in nl43_status") + + if changes_made: + conn.commit() + print("\n✓ Migration completed successfully") + print(" Added wedge-recovery fields to nl43_config and nl43_status") + else: + print("\n✓ All recovery fields already exist - no changes needed") + + conn.close() + return True + + except Exception as e: + print(f"❌ Migration failed: {e}") + return False + + +if __name__ == "__main__": + success = migrate() + sys.exit(0 if success else 1) diff --git a/test_wedge_recovery.py b/test_wedge_recovery.py new file mode 100644 index 0000000..bb3b46b --- /dev/null +++ b/test_wedge_recovery.py @@ -0,0 +1,377 @@ +#!/usr/bin/env python3 +""" +End-to-end test for wedge detection and recovery (app/recovery.py). + +Spins up a fake NL-43 (control + FTP listeners speaking just enough of the +ASCII protocol) that can be "wedged" — control listener killed while FTP +stays up, exactly the failure signature from the Feb 2026 investigation. +A local HTTP server stands in for the relay webhook; hitting it "power +cycles" the fake device. + +Runs against a throwaway SQLite DB in a temp directory — does not touch +data/slmm.db. + +Usage: + python3 test_wedge_recovery.py + +Exit code 0 = all tests passed. +""" + +import asyncio +import os +import sys +import tempfile +import threading +import time +from http.server import BaseHTTPRequestHandler, HTTPServer +from pathlib import Path + +# --- Run from a temp dir so app.database creates a throwaway DB ------------ +REPO_ROOT = Path(__file__).resolve().parent +sys.path.insert(0, str(REPO_ROOT)) +WORKDIR = tempfile.mkdtemp(prefix="slmm-wedge-test-") +os.chdir(WORKDIR) + +# Confirm fast in tests +os.environ.setdefault("WEDGE_CONFIRM_FAILURES", "2") + +from app.database import Base, engine, SessionLocal # noqa: E402 +from app.models import NL43Config, NL43Status # noqa: E402 +from app.recovery import classify_failure, probe_tcp, recovery_manager # noqa: E402 + +Base.metadata.create_all(bind=engine) + +PASS = 0 +FAIL = 0 + + +def check(name: str, condition: bool, detail: str = ""): + global PASS, FAIL + if condition: + PASS += 1 + print(f" ✓ {name}") + else: + FAIL += 1 + print(f" ✗ {name} {('— ' + detail) if detail else ''}") + + +# --------------------------------------------------------------------------- +# Fake NL-43 device +# --------------------------------------------------------------------------- + +class FakeNL43: + """Minimal NL-43: control port speaks the ASCII protocol, FTP port + just accepts. Can be wedged (control listener killed, FTP alive).""" + + def __init__(self): + self.control_server = None + self.ftp_server = None + self.control_port = None + self.ftp_port = None + self.measurement_state = "Stop" + self.index = 1 + self.commands = [] # every command line received + + async def start(self): + self.ftp_server = await asyncio.start_server(self._handle_ftp, "127.0.0.1", 0) + self.ftp_port = self.ftp_server.sockets[0].getsockname()[1] + await self.start_control() + + async def start_control(self): + self.control_server = await asyncio.start_server( + self._handle_control, "127.0.0.1", 0 if self.control_port is None else self.control_port + ) + self.control_port = self.control_server.sockets[0].getsockname()[1] + + async def wedge(self): + """Kill the control listener (existing + new connections die). FTP stays up.""" + if self.control_server: + self.control_server.close() + await self.control_server.wait_closed() + self.control_server = None + + async def power_cycle(self, boot_delay: float = 2.0): + """Simulate a power cycle: everything down, then back up after boot_delay. + Power loss stops any running measurement.""" + await self.wedge() + self.measurement_state = "Stop" + await asyncio.sleep(boot_delay) + await self.start_control() + + async def stop(self): + for srv in (self.control_server, self.ftp_server): + if srv: + srv.close() + await srv.wait_closed() + self.control_server = None + self.ftp_server = None + + async def _handle_ftp(self, reader, writer): + try: + writer.write(b"220 Connection Ready\r\n") + await writer.drain() + await reader.read(1024) + except Exception: + pass + finally: + writer.close() + + async def _handle_control(self, reader, writer): + try: + while True: + line = await reader.readuntil(b"\n") + cmd = line.decode(errors="ignore").strip() + if not cmd: + continue + self.commands.append(cmd) + writer.write(self._respond(cmd)) + await writer.drain() + except (asyncio.IncompleteReadError, ConnectionResetError): + pass + finally: + try: + writer.close() + except Exception: + pass + + def _respond(self, cmd: str) -> bytes: + ok = b"R+0000\r\n" + if cmd == "Measure?": + return ok + f"{self.measurement_state}\r\n".encode() + if cmd == "Measure,Start": + self.measurement_state = "Start" + return ok + if cmd == "Measure,Stop": + self.measurement_state = "Stop" + return ok + if cmd.startswith("Clock,"): + return ok + if cmd == "Store Name?": + return ok + f"{self.index:04d}\r\n".encode() + if cmd.startswith("Store Name,"): + self.index = int(cmd.split(",")[1]) + return ok + if cmd == "Overwrite?": + return ok + b"None\r\n" + if cmd == "DOD?": + return ok + b"1,55.5,54.2,60.1,50.3,72.8\r\n" + if cmd.startswith("Sleep Mode"): + return ok + (b"Off\r\n" if cmd.endswith("?") else b"") + return b"R+0001\r\n" # command error + + +# --------------------------------------------------------------------------- +# Fake relay webhook +# --------------------------------------------------------------------------- + +class FakeRelay: + """HTTP server standing in for a relay controller. A request to /cycle + power-cycles the fake device.""" + + def __init__(self, device: FakeNL43, loop: asyncio.AbstractEventLoop): + self.device = device + self.loop = loop + self.hits = 0 + relay = self + + class Handler(BaseHTTPRequestHandler): + def _handle(self): + relay.hits += 1 + # Schedule the device power cycle on the asyncio loop + asyncio.run_coroutine_threadsafe(relay.device.power_cycle(boot_delay=2.0), relay.loop) + self.send_response(200) + self.end_headers() + self.wfile.write(b"cycling") + + do_GET = _handle + do_POST = _handle + + def log_message(self, *args): + pass + + self.server = HTTPServer(("127.0.0.1", 0), Handler) + self.port = self.server.server_address[1] + self.thread = threading.Thread(target=self.server.serve_forever, daemon=True) + self.thread.start() + + @property + def url(self): + return f"http://127.0.0.1:{self.port}/cycle" + + def stop(self): + self.server.shutdown() + + +# --------------------------------------------------------------------------- +# Tests +# --------------------------------------------------------------------------- + +async def test_classifier(device: FakeNL43): + print("\n[1] Failure classification") + + # Healthy device + c = await classify_failure("127.0.0.1", device.control_port, device.ftp_port) + check("healthy device classifies as ok", c["state"] == "ok", str(c)) + + # Wedged: control refused, FTP alive — the definitive signature + await device.wedge() + c = await classify_failure("127.0.0.1", device.control_port, device.ftp_port) + check("wedged device classifies as wedged", c["state"] == "wedged", str(c)) + check("wedge confidence is high (FTP alive)", c["confidence"] == "high", str(c)) + + # Wedged without FTP probe available + c = await classify_failure("127.0.0.1", device.control_port, None) + check("wedge without FTP probe is medium confidence", + c["state"] == "wedged" and c["confidence"] == "medium", str(c)) + + # Offline: unroutable address times out (TEST-NET-1) + c = await classify_failure("192.0.2.1", 2255, None) + check("unreachable host classifies as offline", c["state"] == "offline", str(c)) + + await device.start_control() + probe = await probe_tcp("127.0.0.1", device.control_port) + check("control port back open after un-wedge", probe == "open", probe) + + +async def test_webhook_recovery(device: FakeNL43, relay: FakeRelay): + print("\n[2] End-to-end webhook recovery (wedge → relay cycle → reconnect → resume)") + + db = SessionLocal() + cfg = NL43Config( + unit_id="TEST-NL43", + host="127.0.0.1", + tcp_port=device.control_port, + ftp_port=device.ftp_port, + tcp_enabled=True, + ftp_enabled=True, + auto_recovery_enabled=True, + reset_backend="webhook", + reset_webhook_url=relay.url, + reset_webhook_method="POST", + recovery_boot_wait_seconds=10, # validator min; device boots in 2s + recovery_reconnect_timeout=60, + recovery_auto_resume=True, + recovery_max_attempts=3, + recovery_window_minutes=360, + ) + db.add(cfg) + # Device was measuring before the wedge + status = NL43Status(unit_id="TEST-NL43", measurement_state="Start", consecutive_failures=2) + db.add(status) + db.commit() + + # Wedge it + device.measurement_state = "Start" + await device.wedge() + + classification = await classify_failure("127.0.0.1", device.control_port, device.ftp_port) + check("pre-recovery classification is wedged", classification["state"] == "wedged") + + started = recovery_manager.start_recovery(cfg, classification) + check("recovery task started", started) + + # Wait for the recovery task to finish (boot wait 10s + commands ≈ 20s) + deadline = time.time() + 90 + while recovery_manager.is_active("TEST-NL43") and time.time() < deadline: + await asyncio.sleep(1) + check("recovery task completed", not recovery_manager.is_active("TEST-NL43")) + + check("relay webhook was hit", relay.hits == 1, f"hits={relay.hits}") + + db.expire_all() + status = db.query(NL43Status).filter_by(unit_id="TEST-NL43").first() + check("recovery_state is idle", status.recovery_state == "idle", status.recovery_state) + check("connection_state is ok", status.connection_state == "ok", status.connection_state) + check("consecutive_failures reset", status.consecutive_failures == 0) + check("last_recovery_result records success", + status.last_recovery_result and "Recovered" in status.last_recovery_result, + str(status.last_recovery_result)) + + check("measurement resumed on device", device.measurement_state == "Start", device.measurement_state) + check("device received Measure,Start", "Measure,Start" in device.commands) + check("start cycle used overwrite protection", "Overwrite?" in device.commands) + + db.close() + + +async def test_gating(device: FakeNL43): + print("\n[3] Recovery gating (confirmation threshold + attempt limit)") + + db = SessionLocal() + cfg = NL43Config( + unit_id="TEST-GATE", + host="127.0.0.1", + tcp_port=device.control_port, + ftp_port=device.ftp_port, + tcp_enabled=True, + auto_recovery_enabled=True, + reset_backend="webhook", + reset_webhook_url="http://127.0.0.1:1/unreachable", # fails fast + recovery_max_attempts=1, + recovery_window_minutes=360, + ) + db.add(cfg) + status = NL43Status(unit_id="TEST-GATE", consecutive_failures=1) + db.add(status) + db.commit() + + classification = {"state": "wedged", "confidence": "high", "tcp_probe": "refused", "ftp_probe": "open"} + + # Below confirmation threshold (1 < 2) — no recovery + started = recovery_manager.maybe_start_recovery(cfg, status, classification, db) + check("below threshold does not trigger", not started) + + # At threshold — triggers (and fails fast on the dead webhook) + status.consecutive_failures = 2 + db.commit() + started = recovery_manager.maybe_start_recovery(cfg, status, classification, db) + check("at threshold triggers recovery", started) + + deadline = time.time() + 30 + while recovery_manager.is_active("TEST-GATE") and time.time() < deadline: + await asyncio.sleep(0.5) + + db.expire_all() + status = db.query(NL43Status).filter_by(unit_id="TEST-GATE").first() + check("failed webhook marks recovery failed", status.recovery_state == "failed", status.recovery_state) + check("failure recorded in result", + status.last_recovery_result and "FAILED" in status.last_recovery_result, + str(status.last_recovery_result)) + + # Attempt limit (max 1 per window) — second attempt blocked + started = recovery_manager.maybe_start_recovery(cfg, status, classification, db) + check("attempt limit blocks repeat recovery", not started) + + # Disabled flag blocks everything + cfg.auto_recovery_enabled = False + db.commit() + started = recovery_manager.maybe_start_recovery(cfg, status, classification, db) + check("auto_recovery_enabled=False blocks recovery", not started) + + db.close() + + +async def main(): + print(f"Wedge recovery test — temp workdir: {WORKDIR}") + + device = FakeNL43() + await device.start() + relay = FakeRelay(device, asyncio.get_running_loop()) + print(f"Fake NL-43: control=127.0.0.1:{device.control_port}, ftp=127.0.0.1:{device.ftp_port}") + print(f"Fake relay webhook: {relay.url}") + + try: + await test_classifier(device) + await test_webhook_recovery(device, relay) + await test_gating(device) + finally: + relay.stop() + await device.stop() + + print(f"\n{'='*50}\nResults: {PASS} passed, {FAIL} failed") + return FAIL == 0 + + +if __name__ == "__main__": + ok = asyncio.run(main()) + sys.exit(0 if ok else 1)