Compare commits
18 Commits
v0.3.0
..
9c43e68534
| Author | SHA1 | Date | |
|---|---|---|---|
| 9c43e68534 | |||
| aa3e088b64 | |||
| 8c17af4849 | |||
| b954eb8c89 | |||
| 0793e7df01 | |||
| 51dd6b682d | |||
| a7983d2958 | |||
| d6dd2e736b | |||
| af86cf713e | |||
| e3f9ca7f5b | |||
| 450509d210 | |||
| fefa9eace8 | |||
| 98a8d357e5 | |||
| 0a7422eceb | |||
| 996b993cb9 | |||
| 01337696b3 | |||
| a302fd15d4 | |||
| af5ecc1a92 |
@@ -1,5 +1,6 @@
|
||||
/manuals/
|
||||
/data/
|
||||
/data-dev/
|
||||
/SLM-stress-test/stress_test_logs/
|
||||
/SLM-stress-test/tcpdump-runs/
|
||||
|
||||
|
||||
@@ -0,0 +1,403 @@
|
||||
# NL-43 + RX55 TCP “Wedge” Investigation (2255 Refusal) — Full Log & Next Steps
|
||||
**Last updated:** 2026-02-18
|
||||
**Owner:** Brian / serversdown
|
||||
**Context:** Terra-View / SLMM / field-deployed Rion NL-43 behind Sierra Wireless RX55
|
||||
|
||||
---
|
||||
|
||||
## 0) What this document is
|
||||
This is a **comprehensive, chronological** record of the debugging we did to isolate a failure where the **NL-43’s TCP control port (2255) eventually stops accepting connections** (“wedges”), while other services (notably FTP/21) remain reachable.
|
||||
|
||||
This is written to be fed back into future troubleshooting, so it intentionally includes the **full reasoning chain, experiments, commands, packet evidence, and conclusions**.
|
||||
|
||||
---
|
||||
|
||||
## 1) Architecture (as tested)
|
||||
### Network path
|
||||
- **Server (SLMM host):** `10.0.0.40`
|
||||
- **RX55 WAN IP:** `63.45.161.30`
|
||||
- **RX55 LAN subnet:** `192.168.1.0/24`
|
||||
- **RX55 LAN gateway:** `192.168.1.1`
|
||||
- **NL-43 LAN IP:** `192.168.1.10` (confirmed via ARP OUI + ping; see LAN validation)
|
||||
|
||||
### RX55 details
|
||||
- **Sierra Wireless RX55**
|
||||
- **OS:** 5.2
|
||||
- **Firmware:** `01.14.24.00`
|
||||
- **Carrier:** Verizon LTE (Band 66)
|
||||
|
||||
### Port forwarding rules (RX55)
|
||||
- **WAN:2255 → NL-43:2255** (NL-43 TCP control)
|
||||
- **WAN:21 → NL-43:21** (NL-43 FTP control)
|
||||
|
||||
You also experimented with additional forwards:
|
||||
- **WAN:2253 → NL-43:2255** (test)
|
||||
- **WAN:2253 → NL-43:2253** (test)
|
||||
- **WAN:4450 → NL-43:4450** (test)
|
||||
|
||||
**Important:** Rule “Input zone / interface” was set to **WAN-NAT**, and Source IP left as **Any IPv4**. This is correct for inbound port-forward behavior on Sierra OS 5.x.
|
||||
|
||||
---
|
||||
|
||||
## 2) Original problem statement (the “wedge”)
|
||||
After running for hours, the NL-43 becomes unreachable over TCP control.
|
||||
|
||||
### Symptom signature (WAN-side)
|
||||
- Client attempts to connect to `63.45.161.30:2255`
|
||||
- Instead of timing out, the client gets **connection refused** quickly.
|
||||
- Packet-level: SYN from client → **RST,ACK** back (meaning active refusal vs silent drop)
|
||||
|
||||
### Critical operational behavior
|
||||
- **Power cycling the NL-43 fixes it.**
|
||||
- **Power cycling the RX55 does NOT fix it.**
|
||||
- FTP sometimes remains available even while TCP control (2255) is dead.
|
||||
|
||||
This combination is what forced us to determine whether:
|
||||
- The RX55 is rejecting connections, OR
|
||||
- The NL-43 is no longer listening on 2255, OR
|
||||
- Something about the RX55 path triggers the NL-43’s control listener to die.
|
||||
|
||||
---
|
||||
|
||||
## 3) Event timeline evidence (SLMM logs)
|
||||
A concrete wedge window was observed on **2026-02-18**:
|
||||
|
||||
- 10:55:46 AM — Poll success (Start)
|
||||
- 11:00:28 AM — Measurement STOPPED (scheduled stop/download cycle succeeded)
|
||||
- 11:55:50 AM — Poll success (Stop)
|
||||
- 12:55:55 PM — Poll success (Stop)
|
||||
- **1:55:58 PM — Poll failed (attempt 1/3): Errno 111 (connection refused)**
|
||||
- 2:56:02 PM — Poll failed (attempt 2/3): Errno 111 (connection refused)
|
||||
|
||||
Key interpretation:
|
||||
- The wedge occurred sometime between **12:55 and 1:55**.
|
||||
- The failure type is **refused**, not timeout.
|
||||
|
||||
---
|
||||
|
||||
## 4) Early hypotheses (before proof)
|
||||
We considered two main buckets:
|
||||
|
||||
### A) NL-43-side failure (most suspicious)
|
||||
- NL-43 TCP control service crashes / exits / unbinds from 2255
|
||||
- socket leak / accept backlog exhaustion
|
||||
- “single control session allowed” and it gets stuck thinking a session is active
|
||||
- mode/service manager bug (service restart fails after other activities)
|
||||
- firmware bug in TCP daemon
|
||||
|
||||
### B) RX55-side failure (possible trigger / less likely once FTP works)
|
||||
- NAT/forwarding table corruption
|
||||
- firewall behavior
|
||||
- helper/ALG interference
|
||||
- MSS/MTU weirdness causing edge-case behavior
|
||||
- session churn behavior causing downstream issues
|
||||
|
||||
---
|
||||
|
||||
## 5) Key experiments and what they proved
|
||||
|
||||
### 5.1) LAN-only stability test (No RX55 path)
|
||||
**Test:** NL-43 tested directly on LAN (no modem path involved).
|
||||
- Ran **24+ hours**
|
||||
- Scheduler start/stop cycles worked
|
||||
- Stress test: **500 commands @ 1/sec** → no failure
|
||||
- Response time trend decreased (not degrading)
|
||||
|
||||
**Result:** The NL-43 appears stable in a “pure LAN” environment.
|
||||
|
||||
**Interpretation:** The trigger is likely related to the RX55/WAN environment, connection patterns, or service switching patterns—not just simple uptime.
|
||||
|
||||
---
|
||||
|
||||
### 5.2) Port-forward behavior: timeout vs refused (RX55 behavior characterization)
|
||||
You observed:
|
||||
|
||||
- **If a WAN port is NOT forwarded (no rule):** connecting to that port **times out** (silent drop)
|
||||
- **If a WAN port IS forwarded to NL-43 but nothing listens:** it **actively refuses** (RST)
|
||||
|
||||
Concrete example:
|
||||
- Port **4450** with no rule → timeout
|
||||
- Port **4450 → NL-43:4450** rule created → connection refused
|
||||
|
||||
**Interpretation:** This confirms the RX55 is actually forwarding packets to the NL-43 when a rule exists. “Refused” is consistent with the NL-43 (or RX55 relay behavior) responding quickly because the packet reached the target.
|
||||
|
||||
Important nuance:
|
||||
- A “refused” on forwarded ports does **not** automatically prove the NL-43 is the one generating RST, because NAT hides the inside host and the RX55 could reject on behalf of an unreachable target. We needed a LAN-side proof test to close the loop.
|
||||
|
||||
---
|
||||
|
||||
### 5.3) UDP test confusion (and resolution)
|
||||
You ran:
|
||||
|
||||
```bash
|
||||
nc -vzu 63.45.161.30 2255
|
||||
nc -vz 63.45.161.30 2255
|
||||
```
|
||||
|
||||
Observed:
|
||||
- UDP: “succeeded”
|
||||
- TCP: “connection refused”
|
||||
|
||||
Resolution:
|
||||
- UDP has **no handshake**. netcat prints “succeeded” if it doesn’t immediately receive an ICMP unreachable. It does **not** mean a UDP service exists.
|
||||
- TCP refused is meaningful: a RST implies “no listener” or “actively rejected.”
|
||||
|
||||
**Net effect:** UDP test did not change the diagnosis.
|
||||
|
||||
---
|
||||
|
||||
### 5.4) Packet capture proof (WAN-side)
|
||||
You captured a Wireshark/tcpdump summary with these key patterns:
|
||||
|
||||
#### Port 2255 (TCP control)
|
||||
Example:
|
||||
- `10.0.0.40 → 63.45.161.30:2255` SYN
|
||||
- `63.45.161.30 → 10.0.0.40` **RST, ACK** within ~50ms
|
||||
|
||||
This happened repeatedly.
|
||||
|
||||
#### Port 2253 (test port)
|
||||
Multiple SYN attempts to 2253 showed **retransmissions and no response**, i.e., **silent drop** (consistent with no rule or not forwarded at that moment).
|
||||
|
||||
#### Port 21 (FTP)
|
||||
Clean 3-way handshake:
|
||||
- SYN → SYN/ACK → ACK
|
||||
Then:
|
||||
- FTP server banner: `220 Connection Ready`
|
||||
Then:
|
||||
- `530 Not logged in` (because SLMM was sending non-FTP “requests” as an experiment)
|
||||
Session closes cleanly.
|
||||
|
||||
**Key takeaway from capture:**
|
||||
- TCP transport to NL-43 via RX55 is definitely working (port 21 proves it).
|
||||
- Port 2255 is being actively refused.
|
||||
|
||||
This strongly suggested “2255 listener is gone,” but still didn’t fully prove whether the refusal was generated internally by NL-43 or by RX55 on behalf of NL-43.
|
||||
|
||||
---
|
||||
|
||||
## 6) The decisive experiment: LAN-side test while wedged (final proof)
|
||||
Because the RX55 does not offer SSH, the plan was to test from **inside the LAN behind the RX55**.
|
||||
|
||||
### 6.1) Physical LAN tap setup
|
||||
Constraint:
|
||||
- NL-43 has only one Ethernet port.
|
||||
|
||||
Solution:
|
||||
- Insert an unmanaged switch:
|
||||
- RX55 LAN → switch
|
||||
- NL-43 → switch
|
||||
- Windows 10 laptop → switch
|
||||
|
||||
This creates a shared L2 segment where the laptop can test NL-43 directly.
|
||||
|
||||
### 6.2) Windows LAN validation
|
||||
On the Windows laptop:
|
||||
|
||||
- `ipconfig` showed:
|
||||
- IP: `192.168.1.100`
|
||||
- Gateway: `192.168.1.1` (RX55)
|
||||
- Initial `arp -a` only showed RX55, not NL-43.
|
||||
|
||||
You then:
|
||||
- pinged likely host addresses and discovered NL-43 responds on **192.168.1.10**
|
||||
- `arp -a` then showed:
|
||||
- `192.168.1.10 → 00-10-50-14-0a-d8`
|
||||
- OUI `00-10-50` recognized as **Rion** (matches NL-43)
|
||||
|
||||
So LAN identities were confirmed:
|
||||
- RX55: `192.168.1.1`
|
||||
- NL-43: `192.168.1.10`
|
||||
|
||||
### 6.3) The LAN port tests (the smoking gun)
|
||||
From Windows:
|
||||
|
||||
```powershell
|
||||
Test-NetConnection -ComputerName 192.168.1.10 -Port 2255
|
||||
Test-NetConnection -ComputerName 192.168.1.10 -Port 21
|
||||
```
|
||||
|
||||
Results (while the unit was “wedged” from the WAN perspective):
|
||||
- **2255:** `TcpTestSucceeded : False`
|
||||
- **21:** `TcpTestSucceeded : True`
|
||||
|
||||
**Conclusion (PROVEN):**
|
||||
- The NL-43 is reachable on the LAN
|
||||
- FTP port 21 is alive
|
||||
- **The NL-43 is NOT listening on TCP port 2255**
|
||||
- Therefore the RX55 is not the root cause of the refusal. The WAN refusal is consistent with the NL-43 having no listener on 2255.
|
||||
|
||||
This is now settled.
|
||||
|
||||
---
|
||||
|
||||
## 7) What we learned (final conclusions)
|
||||
### 7.1) RX55 innocence (for this failure mode)
|
||||
The RX55 is not “randomly rejecting” or “breaking TCP” in the way originally feared.
|
||||
|
||||
It successfully forwards and supports TCP to the NL-43 on port 21, and the LAN-side test proves the 2255 failure exists *even without NAT/WAN involvement*.
|
||||
|
||||
### 7.2) NL-43 control listener failure
|
||||
The NL-43’s TCP control service (port 2255) stops listening while:
|
||||
- the device remains alive
|
||||
- the LAN stack remains alive (ping)
|
||||
- FTP remains alive (port 21)
|
||||
|
||||
This looks like one of:
|
||||
- control daemon crash/exit
|
||||
- service unbind
|
||||
- stuck service state (e.g., “busy” / “session active forever”)
|
||||
- resource leak (sockets/file descriptors) specific to the control service
|
||||
- firmware service manager bug (start/stop of services fails after certain sequences)
|
||||
|
||||
---
|
||||
|
||||
## 8) Additional constraint discovered: “Web App mode” conflicts
|
||||
You noted an important operational constraint:
|
||||
|
||||
> Turning on the web app disables other interfaces like TCP and FTP.
|
||||
|
||||
Meaning the NL-43 appears to have mutually exclusive service/mode behavior (or at least serious conflicts). That matters because:
|
||||
- If any workflow toggles modes (explicitly or implicitly), it could destabilize the service lifecycle.
|
||||
- It reduces the possibility of using “web UI toggle” as an easy remote recovery mechanism **if** it disables the services needed.
|
||||
|
||||
We have not yet run a controlled long test to determine whether:
|
||||
- mode switching contributes directly to the 2255 listener dying, OR
|
||||
- it happens even in a pure TCP-only mode with no switching.
|
||||
|
||||
---
|
||||
|
||||
## 9) Immediate operational decision (field tomorrow)
|
||||
Because the device is needed in the field immediately, you chose:
|
||||
- **Old-school manual deployment**
|
||||
- **Manual SD card downloads**
|
||||
- Avoid reliance on 2255/TCP control and remote workflows for now.
|
||||
|
||||
**Important operational note:**
|
||||
The 2255 listener dying does not necessarily stop the NL-43 from measuring; it primarily breaks remote control/polling. Manual SD workflow sidesteps the entire remote control dependency.
|
||||
|
||||
---
|
||||
|
||||
## 10) What’s next (future work — when the unit is back)
|
||||
Because long tests can’t be run before tomorrow, the plan is to resume in a few weeks with controlled experiments designed to isolate the trigger and develop an operational mitigation.
|
||||
|
||||
### 10.1) Controlled experiment matrix (recommended)
|
||||
Run each test for 24–72 hours, or until wedge occurs, and record:
|
||||
- number of TCP connects
|
||||
- whether connections are persistent
|
||||
- whether FTP is used
|
||||
- whether any mode toggling is performed
|
||||
- time-to-wedge
|
||||
|
||||
#### Test A — TCP-only (ideal baseline)
|
||||
- TCP control only (2255)
|
||||
- **True persistent connection** (open once, keep forever)
|
||||
- No FTP
|
||||
- No web mode toggling
|
||||
|
||||
Outcome interpretation:
|
||||
- If stable: connection churn and/or FTP/mode switching is the trigger.
|
||||
- If wedges anyway: pure 2255 daemon leak/bug.
|
||||
|
||||
#### Test B — TCP with connection churn
|
||||
- Same as A but intentionally reconnect on a schedule (current SLMM behavior)
|
||||
- No FTP
|
||||
|
||||
Outcome:
|
||||
- If this wedges but A doesn’t: churn is the trigger.
|
||||
|
||||
#### Test C — FTP activity + TCP
|
||||
- Introduce scheduled FTP sessions (downloads) while using TCP control
|
||||
- Observe whether wedge correlates with FTP use or with post-download periods.
|
||||
|
||||
Outcome:
|
||||
- If wedge correlates with FTP, suspect internal service lifecycle conflict.
|
||||
|
||||
#### Test D — Web mode interaction (only if safe/possible)
|
||||
- Evaluate what toggling web mode does to TCP/FTP services.
|
||||
- Determine if any remote-safe “soft reset” exists.
|
||||
|
||||
---
|
||||
|
||||
## 11) Mitigation options (ranked)
|
||||
### Option 1 — Make SLMM truly persistent (highest probability of success)
|
||||
If the NL-43 wedges due to session churn or leaked socket states, the best mitigation is:
|
||||
- Open one TCP socket per device
|
||||
- Keep it open indefinitely
|
||||
- Use OS keepalive
|
||||
- Do **not** rotate connections on timers
|
||||
- Reconnect only when the socket actually dies
|
||||
|
||||
This reduces:
|
||||
- connect/close cycles
|
||||
- NAT edge-case exposure
|
||||
- resource churn inside NL-43
|
||||
|
||||
### Option 2 — Service “soft reset” (if possible without disabling required services)
|
||||
If there exists any way to restart the 2255 service without power cycling:
|
||||
- LAN TCP toggle (if it doesn’t require web mode)
|
||||
- any “restart comms” command (unknown)
|
||||
- any maintenance menu sequence
|
||||
then SLMM could:
|
||||
- detect wedge
|
||||
- trigger soft reset
|
||||
- recover automatically
|
||||
|
||||
Current constraint: web app mode appears to disable other services, so this may not be viable.
|
||||
|
||||
### Option 3 — Hardware watchdog power cycle (industrial but reliable)
|
||||
If this is a firmware bug with no clean workaround:
|
||||
- Add a remotely controlled relay/power switch
|
||||
- On wedge detection, power-cycle NL-43 automatically
|
||||
- Optionally schedule a nightly power cycle to prevent leak accumulation
|
||||
|
||||
This is “field reality” and often the only long-term move with embedded devices.
|
||||
|
||||
### Option 4 — Vendor escalation (Rion)
|
||||
You now have excellent evidence:
|
||||
- LAN-side proof: 2255 dead while 21 alive
|
||||
- WAN packet evidence
|
||||
- clear isolation of RX55 innocence
|
||||
|
||||
This is strong enough to send to Rion support as a firmware defect report.
|
||||
|
||||
---
|
||||
|
||||
## 12) Repro “wedge bundle” checklist (for future captures)
|
||||
When the wedge happens again, capture these before power cycling:
|
||||
|
||||
1) From server:
|
||||
- `nc -vz 63.45.161.30 2255` (expect refused)
|
||||
- `nc -vz 63.45.161.30 21` (expect success if FTP alive)
|
||||
|
||||
2) From LAN side (via switch/laptop):
|
||||
- `Test-NetConnection 192.168.1.10 -Port 2255`
|
||||
- `Test-NetConnection 192.168.1.10 -Port 21`
|
||||
|
||||
3) Optional: packet capture around the refused attempt.
|
||||
|
||||
4) Record:
|
||||
- last successful poll timestamp
|
||||
- last FTP session timestamp
|
||||
- any scheduled start/stop/download cycles near wedge time
|
||||
- SLMM connection reuse/rotation settings in effect
|
||||
|
||||
---
|
||||
|
||||
## 13) Final, current-state summary (as of 2026-02-18)
|
||||
- The issue is **NOT** the RX55 rejecting inbound connections.
|
||||
- The NL-43 is **alive**, reachable on LAN, and FTP works.
|
||||
- The NL-43’s **TCP control listener on 2255 stops listening** while the device remains otherwise healthy.
|
||||
- The wedge can occur hours after successful operations.
|
||||
- The unit is needed in the field immediately, so investigation pauses.
|
||||
- Next phase: controlled tests to isolate trigger + implement mitigation (persistent socket or watchdog reset).
|
||||
|
||||
---
|
||||
|
||||
## 14) Notes / misc observations
|
||||
- The Wireshark trace showed repeated FTP sessions were opened and closed cleanly, but SLMM’s “FTP requests” were not valid FTP (causing `530 Not logged in`). That was part of experimentation, not a normal workflow.
|
||||
- UDP “success” via netcat is not meaningful because UDP has no handshake; it simply indicates no ICMP unreachable was returned.
|
||||
|
||||
---
|
||||
|
||||
**End of document.**
|
||||
+244
@@ -0,0 +1,244 @@
|
||||
"""
|
||||
Threshold alert engine.
|
||||
|
||||
Each unit can have any number of AlertRules. A rule is evaluated against the
|
||||
unit's live monitor snapshots via a small per-(unit, rule) state machine:
|
||||
|
||||
IDLE --(metric exceeds threshold for duration_s)--> ACTIVE (fire ONSET)
|
||||
ACTIVE --(metric recovers past hysteresis for duration_s)--> IDLE (fire CLEAR)
|
||||
|
||||
duration_s debounces both edges; clear_margin_db adds hysteresis so a level
|
||||
hovering at the threshold doesn't flap. Onset and clear are distinct events.
|
||||
|
||||
The state-machine logic (`_evaluate_step`) is intentionally pure — no DB, no
|
||||
real clock — so it can be unit-tested with a synthetic level series and a fake
|
||||
clock. The AlertEvaluator wraps it with rule loading, scheduling, persistence,
|
||||
and dispatch. Dispatch is a server log for now (POC); the seam to POST events to
|
||||
a Terra-View webhook (email/SMS) is _dispatch().
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Local timezone offset for schedule windows (same env var services.py uses).
|
||||
_TZ_OFFSET_HOURS = float(os.getenv("TIMEZONE_OFFSET", "-5"))
|
||||
|
||||
# How long to cache a unit's rules before re-querying the DB (rules change rarely).
|
||||
_RULE_CACHE_TTL_S = 15.0
|
||||
|
||||
|
||||
@dataclass
|
||||
class RuleState:
|
||||
"""In-memory runtime state for one (unit, rule)."""
|
||||
phase: str = "idle" # "idle" | "active"
|
||||
edge_since: Optional[float] = None # when the current edge condition began (clock time)
|
||||
peak: float = 0.0
|
||||
event_id: Optional[int] = None # the open AlertEvent row (for the clear update)
|
||||
|
||||
|
||||
def _exceeds(value: float, rule) -> bool:
|
||||
if rule.comparison == "below":
|
||||
return value < rule.threshold_db
|
||||
return value > rule.threshold_db
|
||||
|
||||
|
||||
def _recovered(value: float, rule) -> bool:
|
||||
margin = rule.clear_margin_db or 0.0
|
||||
if rule.comparison == "below":
|
||||
return value > rule.threshold_db + margin
|
||||
return value < rule.threshold_db - margin
|
||||
|
||||
|
||||
def _evaluate_step(state: RuleState, value: float, now: float, rule) -> Optional[str]:
|
||||
"""Advance the state machine by one reading.
|
||||
|
||||
Pure: mutates `state`, returns 'onset' | 'clear' | None. `now` is injected so
|
||||
tests can drive a fake clock.
|
||||
"""
|
||||
duration = rule.duration_s or 0
|
||||
|
||||
if state.phase == "idle":
|
||||
if _exceeds(value, rule):
|
||||
if state.edge_since is None:
|
||||
state.edge_since = now
|
||||
if now - state.edge_since >= duration:
|
||||
state.phase = "active"
|
||||
state.edge_since = None
|
||||
state.peak = value
|
||||
return "onset"
|
||||
else:
|
||||
state.edge_since = None
|
||||
return None
|
||||
|
||||
# active
|
||||
if rule.comparison == "below":
|
||||
state.peak = min(state.peak, value)
|
||||
else:
|
||||
state.peak = max(state.peak, value)
|
||||
|
||||
if _recovered(value, rule):
|
||||
if state.edge_since is None:
|
||||
state.edge_since = now
|
||||
if now - state.edge_since >= duration:
|
||||
state.phase = "idle"
|
||||
state.edge_since = None
|
||||
return "clear"
|
||||
else:
|
||||
state.edge_since = None
|
||||
return None
|
||||
|
||||
|
||||
def _in_window(now_minutes: int, start: str, end: str) -> bool:
|
||||
"""Is now_minutes (minutes since local midnight) within [start, end)?
|
||||
Handles wraparound windows like 22:00–07:00."""
|
||||
def _m(s: str) -> int:
|
||||
h, m = s.split(":")
|
||||
return int(h) * 60 + int(m)
|
||||
s, e = _m(start), _m(end)
|
||||
if s == e:
|
||||
return True
|
||||
if s < e:
|
||||
return s <= now_minutes < e
|
||||
return now_minutes >= s or now_minutes < e # wraparound
|
||||
|
||||
|
||||
class AlertEvaluator:
|
||||
def __init__(self):
|
||||
self._states: Dict[Tuple[str, int], RuleState] = {}
|
||||
self._rule_cache: Dict[str, Tuple[float, list]] = {} # unit_id -> (fetched_at, rules)
|
||||
logger.info("[ALERT] rule-based evaluator ready")
|
||||
|
||||
async def evaluate(self, unit_id: str, snap) -> None:
|
||||
"""Evaluate every enabled rule for this unit against one snapshot."""
|
||||
rules = self._get_rules(unit_id)
|
||||
if not rules:
|
||||
return
|
||||
now = asyncio.get_running_loop().time()
|
||||
for rule in rules:
|
||||
if not self._in_schedule(rule):
|
||||
continue
|
||||
raw = getattr(snap, rule.metric, None)
|
||||
try:
|
||||
value = float(raw)
|
||||
except (TypeError, ValueError):
|
||||
continue # missing / non-numeric ("-.-")
|
||||
state = self._states.setdefault((unit_id, rule.id), RuleState())
|
||||
action = _evaluate_step(state, value, now, rule)
|
||||
if action == "onset":
|
||||
await self._on_onset(unit_id, rule, value, state)
|
||||
elif action == "clear":
|
||||
await self._on_clear(unit_id, rule, value, state)
|
||||
|
||||
# -- rule loading (cached) ----------------------------------------------
|
||||
|
||||
def _get_rules(self, unit_id: str) -> list:
|
||||
loop_now = asyncio.get_running_loop().time()
|
||||
cached = self._rule_cache.get(unit_id)
|
||||
if cached and loop_now - cached[0] < _RULE_CACHE_TTL_S:
|
||||
return cached[1]
|
||||
rules = self._load_rules(unit_id)
|
||||
self._rule_cache[unit_id] = (loop_now, rules)
|
||||
return rules
|
||||
|
||||
def _load_rules(self, unit_id: str) -> list:
|
||||
from app.database import SessionLocal
|
||||
from app.models import AlertRule
|
||||
db = SessionLocal()
|
||||
try:
|
||||
return db.query(AlertRule).filter_by(unit_id=unit_id, enabled=True).all()
|
||||
except Exception as e:
|
||||
logger.warning(f"[ALERT] failed to load rules for {unit_id}: {e}")
|
||||
return []
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def invalidate(self, unit_id: Optional[str] = None) -> None:
|
||||
"""Drop cached rules so a change is picked up immediately."""
|
||||
if unit_id is None:
|
||||
self._rule_cache.clear()
|
||||
else:
|
||||
self._rule_cache.pop(unit_id, None)
|
||||
|
||||
# -- scheduling ----------------------------------------------------------
|
||||
|
||||
def _in_schedule(self, rule) -> bool:
|
||||
if not rule.schedule_start or not rule.schedule_end:
|
||||
day_ok = self._day_ok(rule)
|
||||
return day_ok
|
||||
local = datetime.utcnow() + timedelta(hours=_TZ_OFFSET_HOURS)
|
||||
if not self._day_ok(rule, local):
|
||||
return False
|
||||
return _in_window(local.hour * 60 + local.minute, rule.schedule_start, rule.schedule_end)
|
||||
|
||||
@staticmethod
|
||||
def _day_ok(rule, local: Optional[datetime] = None) -> bool:
|
||||
if not rule.schedule_days:
|
||||
return True
|
||||
if local is None:
|
||||
local = datetime.utcnow() + timedelta(hours=_TZ_OFFSET_HOURS)
|
||||
allowed = {int(d) for d in str(rule.schedule_days).split(",") if d.strip() != ""}
|
||||
return local.weekday() in allowed # Mon=0
|
||||
|
||||
# -- event persistence + dispatch ---------------------------------------
|
||||
|
||||
async def _on_onset(self, unit_id: str, rule, value: float, state: RuleState) -> None:
|
||||
from app.database import SessionLocal
|
||||
from app.models import AlertEvent
|
||||
db = SessionLocal()
|
||||
try:
|
||||
evt = AlertEvent(
|
||||
rule_id=rule.id, unit_id=unit_id, rule_name=rule.name,
|
||||
metric=rule.metric, threshold_db=rule.threshold_db,
|
||||
onset_value=value, peak_value=value, status="active",
|
||||
)
|
||||
db.add(evt)
|
||||
db.commit()
|
||||
db.refresh(evt)
|
||||
state.event_id = evt.id
|
||||
except Exception as e:
|
||||
logger.warning(f"[ALERT] failed to record onset for {unit_id}: {e}")
|
||||
finally:
|
||||
db.close()
|
||||
await self._dispatch(
|
||||
"ONSET", unit_id, rule,
|
||||
f"{rule.metric.upper()}={value:.1f} dB "
|
||||
f"{'<' if rule.comparison == 'below' else '>'} {rule.threshold_db:.1f} dB"
|
||||
f"{f' for {rule.duration_s}s' if rule.duration_s else ''}",
|
||||
)
|
||||
|
||||
async def _on_clear(self, unit_id: str, rule, value: float, state: RuleState) -> None:
|
||||
peak = state.peak
|
||||
from app.database import SessionLocal
|
||||
from app.models import AlertEvent
|
||||
db = SessionLocal()
|
||||
try:
|
||||
if state.event_id is not None:
|
||||
evt = db.query(AlertEvent).filter_by(id=state.event_id).first()
|
||||
if evt:
|
||||
evt.clear_at = datetime.utcnow()
|
||||
evt.peak_value = peak
|
||||
evt.status = "cleared"
|
||||
db.commit()
|
||||
except Exception as e:
|
||||
logger.warning(f"[ALERT] failed to record clear for {unit_id}: {e}")
|
||||
finally:
|
||||
db.close()
|
||||
state.event_id = None
|
||||
await self._dispatch(
|
||||
"CLEAR", unit_id, rule,
|
||||
f"recovered to {value:.1f} dB (peak {peak:.1f} dB)",
|
||||
)
|
||||
|
||||
async def _dispatch(self, kind: str, unit_id: str, rule, detail: str) -> None:
|
||||
"""POC dispatch: server log. Swap in a Terra-View webhook (email/SMS) here."""
|
||||
logger.warning(f"[ALERT:{kind}] {unit_id} '{rule.name}': {detail}")
|
||||
|
||||
|
||||
# Module-level singleton (the monitor calls alert_evaluator.evaluate per snapshot)
|
||||
alert_evaluator = AlertEvaluator()
|
||||
@@ -8,6 +8,7 @@ for fast API access without querying devices on every request.
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Optional
|
||||
|
||||
@@ -20,6 +21,11 @@ from app.device_logger import log_device_event, cleanup_old_logs
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Global polling default. Set SLMM_POLLING_ENABLED=false to start an instance in
|
||||
# standby (running but not polling and not holding device connections) — e.g. a
|
||||
# dev box that must not latch onto a device that a prod instance owns.
|
||||
POLLING_ENABLED_DEFAULT = os.getenv("SLMM_POLLING_ENABLED", "true").lower() == "true"
|
||||
|
||||
|
||||
class BackgroundPoller:
|
||||
"""
|
||||
@@ -38,6 +44,8 @@ class BackgroundPoller:
|
||||
self._running = False
|
||||
self._logger = logger
|
||||
self._last_cleanup = None # Track last log cleanup time
|
||||
self._last_pool_log = None # Track last connection pool heartbeat log
|
||||
self._active = POLLING_ENABLED_DEFAULT # Global polling on/off (standby toggle)
|
||||
|
||||
async def start(self):
|
||||
"""Start the background polling task."""
|
||||
@@ -70,15 +78,48 @@ class BackgroundPoller:
|
||||
|
||||
self._logger.info("Background poller stopped")
|
||||
|
||||
def is_active(self) -> bool:
|
||||
"""Whether background polling is currently active (vs standby)."""
|
||||
return self._active
|
||||
|
||||
async def set_active(self, active: bool):
|
||||
"""Globally enable/disable polling at runtime.
|
||||
|
||||
When deactivated, the loop stays alive but polls nothing and releases all
|
||||
device connections, so this SLMM instance stops occupying the devices'
|
||||
single connection slots (e.g. so a prod instance can take over). Runtime
|
||||
state only — on restart the instance returns to SLMM_POLLING_ENABLED.
|
||||
"""
|
||||
self._active = active
|
||||
if active:
|
||||
self._logger.info("[SYSTEM] Background polling ACTIVATED")
|
||||
else:
|
||||
self._logger.info("[SYSTEM] Background polling DEACTIVATED (standby) — releasing connections")
|
||||
await self._release_all_connections()
|
||||
|
||||
async def _release_all_connections(self):
|
||||
"""Gracefully close every pooled device connection (no-op if none)."""
|
||||
from app.services import _connection_pool
|
||||
for device_key in list(_connection_pool.get_stats().get("connections", {})):
|
||||
await _connection_pool.discard(device_key)
|
||||
|
||||
async def _poll_loop(self):
|
||||
"""Main polling loop that runs continuously."""
|
||||
self._logger.info("Background polling loop started")
|
||||
|
||||
while self._running:
|
||||
try:
|
||||
await self._poll_all_devices()
|
||||
except Exception as e:
|
||||
self._logger.error(f"Error in poll loop: {e}", exc_info=True)
|
||||
if self._active:
|
||||
try:
|
||||
await self._poll_all_devices()
|
||||
except Exception as e:
|
||||
self._logger.error(f"Error in poll loop: {e}", exc_info=True)
|
||||
else:
|
||||
# Standby: poll nothing, and keep holding no device connection slots
|
||||
# so another SLMM instance (e.g. prod) can talk to the devices.
|
||||
try:
|
||||
await self._release_all_connections()
|
||||
except Exception as e:
|
||||
self._logger.warning(f"Standby connection release failed: {e}")
|
||||
|
||||
# Run log cleanup once per hour
|
||||
try:
|
||||
@@ -89,6 +130,24 @@ class BackgroundPoller:
|
||||
except Exception as e:
|
||||
self._logger.warning(f"Log cleanup failed: {e}")
|
||||
|
||||
# Log connection pool status every 15 minutes
|
||||
try:
|
||||
now = datetime.utcnow()
|
||||
if self._last_pool_log is None or (now - self._last_pool_log).total_seconds() > 900:
|
||||
from app.services import _connection_pool
|
||||
stats = _connection_pool.get_stats()
|
||||
conns = stats.get("connections", {})
|
||||
if conns:
|
||||
for key, c in conns.items():
|
||||
self._logger.info(
|
||||
f"[POOL] {key} — age={c['age_seconds']}s idle={c['idle_seconds']}s alive={c['alive']}"
|
||||
)
|
||||
else:
|
||||
self._logger.info("[POOL] No active connections in pool")
|
||||
self._last_pool_log = now
|
||||
except Exception as e:
|
||||
self._logger.warning(f"Pool status log failed: {e}")
|
||||
|
||||
# Calculate dynamic sleep interval
|
||||
sleep_time = self._calculate_sleep_interval()
|
||||
self._logger.debug(f"Sleeping for {sleep_time} seconds until next poll cycle")
|
||||
|
||||
+2
-2
@@ -76,12 +76,12 @@ app.include_router(routers.router)
|
||||
|
||||
@app.get("/", response_class=HTMLResponse)
|
||||
def index(request: Request):
|
||||
return templates.TemplateResponse("index.html", {"request": request})
|
||||
return templates.TemplateResponse(request, "index.html")
|
||||
|
||||
|
||||
@app.get("/roster", response_class=HTMLResponse)
|
||||
def roster(request: Request):
|
||||
return templates.TemplateResponse("roster.html", {"request": request})
|
||||
return templates.TemplateResponse(request, "roster.html")
|
||||
|
||||
|
||||
@app.get("/health")
|
||||
|
||||
+53
-1
@@ -1,4 +1,4 @@
|
||||
from sqlalchemy import Column, String, DateTime, Boolean, Integer, Text, func
|
||||
from sqlalchemy import Column, String, DateTime, Boolean, Integer, Float, Text, func
|
||||
from app.database import Base
|
||||
|
||||
|
||||
@@ -41,6 +41,8 @@ class NL43Status(Base):
|
||||
lmax = Column(String, nullable=True) # Maximum level
|
||||
lmin = Column(String, nullable=True) # Minimum level
|
||||
lpeak = Column(String, nullable=True) # Peak level
|
||||
ln1 = Column(String, nullable=True) # Percentile slot LN1 (configurable; device default L5, contract L1)
|
||||
ln2 = Column(String, nullable=True) # Percentile slot LN2 (configurable; device default L10)
|
||||
battery_level = Column(String, nullable=True)
|
||||
power_source = Column(String, nullable=True)
|
||||
sd_remaining_mb = Column(String, nullable=True)
|
||||
@@ -72,3 +74,53 @@ class DeviceLog(Base):
|
||||
level = Column(String, default="INFO") # DEBUG, INFO, WARNING, ERROR
|
||||
category = Column(String, default="GENERAL") # TCP, FTP, POLL, COMMAND, STATE, SYNC
|
||||
message = Column(Text, nullable=False)
|
||||
|
||||
|
||||
class AlertRule(Base):
|
||||
"""A threshold-alert rule evaluated against a unit's live monitor feed.
|
||||
|
||||
Source-agnostic: today it runs over the DOD monitor; the same rule transfers
|
||||
unchanged if a unit's feed is later sourced from FTP intervals.
|
||||
"""
|
||||
|
||||
__tablename__ = "alert_rules"
|
||||
|
||||
id = Column(Integer, primary_key=True, autoincrement=True)
|
||||
unit_id = Column(String, index=True, nullable=False)
|
||||
name = Column(String, nullable=False, default="Alert")
|
||||
metric = Column(String, nullable=False, default="lp") # lp/leq/lmax/lmin/lpeak/ln1/ln2
|
||||
comparison = Column(String, nullable=False, default="above") # above | below
|
||||
threshold_db = Column(Float, nullable=False)
|
||||
duration_s = Column(Integer, nullable=False, default=0) # sustained seconds (0 = instant)
|
||||
clear_margin_db = Column(Float, nullable=False, default=2.0) # hysteresis band
|
||||
cooldown_s = Column(Integer, nullable=False, default=300) # min seconds between onsets
|
||||
# Optional time-of-day scoping (local time). schedule_start/end as "HH:MM";
|
||||
# null = always active. schedule_days = CSV of 0-6 (Mon=0); null = every day.
|
||||
schedule_start = Column(String, nullable=True)
|
||||
schedule_end = Column(String, nullable=True)
|
||||
schedule_days = Column(String, nullable=True)
|
||||
channels = Column(String, nullable=False, default="log") # CSV: log,email,sms
|
||||
recipients = Column(Text, nullable=True) # CSV of emails/phones
|
||||
enabled = Column(Boolean, default=True)
|
||||
created_at = Column(DateTime, default=func.now())
|
||||
|
||||
|
||||
class AlertEvent(Base):
|
||||
"""A fired alert (onset → clear), for history / inbox / acknowledgement."""
|
||||
|
||||
__tablename__ = "alert_events"
|
||||
|
||||
id = Column(Integer, primary_key=True, autoincrement=True)
|
||||
rule_id = Column(Integer, index=True, nullable=False)
|
||||
unit_id = Column(String, index=True, nullable=False)
|
||||
rule_name = Column(String, nullable=True)
|
||||
metric = Column(String, nullable=False)
|
||||
threshold_db = Column(Float, nullable=False)
|
||||
onset_at = Column(DateTime, default=func.now(), index=True)
|
||||
onset_value = Column(Float, nullable=True)
|
||||
peak_value = Column(Float, nullable=True)
|
||||
clear_at = Column(DateTime, nullable=True)
|
||||
status = Column(String, default="active") # active | cleared
|
||||
acknowledged_at = Column(DateTime, nullable=True)
|
||||
acknowledged_by = Column(String, nullable=True)
|
||||
notes = Column(Text, nullable=True)
|
||||
|
||||
+176
@@ -0,0 +1,176 @@
|
||||
"""
|
||||
Per-device live monitor (fan-out hub).
|
||||
|
||||
ONE DOD poll loop per device, broadcast to many subscribers:
|
||||
- browser WebSocket clients (live view) — they no longer each open their own
|
||||
device stream, so the NL43's single-connection limit stops causing the
|
||||
"second viewer sees nothing" contention.
|
||||
- the alert evaluator (threshold alerts), which can keep a device's feed running
|
||||
even with no browser attached.
|
||||
- persistence (each snapshot is written to NL43Status, like the poller does).
|
||||
|
||||
The device's one TCP connection is respected: every poll goes through the same
|
||||
per-device lock + connection pool in services.py, so the monitor, the background
|
||||
poller, and on-demand commands all serialize safely.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
from datetime import datetime
|
||||
from typing import Dict, Optional, Set
|
||||
|
||||
from app.database import SessionLocal
|
||||
from app.models import NL43Config, NL43Status
|
||||
from app.services import NL43Client, persist_snapshot
|
||||
from app.alerts import alert_evaluator
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Sleep between DOD polls. Note the 1s device rate-limit (and DOD?+Measure? per
|
||||
# poll) already paces the effective rate to a few seconds; this is the extra idle.
|
||||
MONITOR_POLL_INTERVAL = float(os.getenv("MONITOR_POLL_INTERVAL", "1.0"))
|
||||
|
||||
|
||||
def _snapshot_payload(snap, unit_id: str, measurement_start_time) -> dict:
|
||||
"""Build the broadcast payload — same shape as the DRD stream, but DOD-sourced
|
||||
so it carries ln1/ln2 (which DRD cannot)."""
|
||||
return {
|
||||
"unit_id": unit_id,
|
||||
"timestamp": datetime.utcnow().isoformat(),
|
||||
"measurement_state": snap.measurement_state,
|
||||
"measurement_start_time": measurement_start_time,
|
||||
"counter": snap.counter,
|
||||
"lp": snap.lp,
|
||||
"leq": snap.leq,
|
||||
"lmax": snap.lmax,
|
||||
"lmin": snap.lmin,
|
||||
"lpeak": snap.lpeak,
|
||||
"ln1": snap.ln1,
|
||||
"ln2": snap.ln2,
|
||||
"raw_payload": snap.raw_payload,
|
||||
}
|
||||
|
||||
|
||||
class DeviceMonitor:
|
||||
"""Owns a single DOD poll loop for one device and fans each snapshot out to
|
||||
all subscribers. Runs while it has at least one browser subscriber OR the
|
||||
server-side keep-alive (alerting) flag is set."""
|
||||
|
||||
def __init__(self, unit_id: str):
|
||||
self.unit_id = unit_id
|
||||
self._subscribers: Set[asyncio.Queue] = set()
|
||||
self._keepalive = False
|
||||
self._task: Optional[asyncio.Task] = None
|
||||
self._lock = asyncio.Lock()
|
||||
|
||||
@property
|
||||
def running(self) -> bool:
|
||||
return self._task is not None and not self._task.done()
|
||||
|
||||
def subscriber_count(self) -> int:
|
||||
return len(self._subscribers)
|
||||
|
||||
def _has_demand(self) -> bool:
|
||||
return bool(self._subscribers) or self._keepalive
|
||||
|
||||
def _ensure_task(self) -> None:
|
||||
if self._task is None or self._task.done():
|
||||
self._task = asyncio.create_task(self._run())
|
||||
|
||||
async def subscribe(self) -> asyncio.Queue:
|
||||
q: asyncio.Queue = asyncio.Queue(maxsize=5)
|
||||
async with self._lock:
|
||||
self._subscribers.add(q)
|
||||
self._ensure_task()
|
||||
return q
|
||||
|
||||
async def unsubscribe(self, q: asyncio.Queue) -> None:
|
||||
async with self._lock:
|
||||
self._subscribers.discard(q)
|
||||
|
||||
async def set_keepalive(self, on: bool) -> None:
|
||||
async with self._lock:
|
||||
self._keepalive = on
|
||||
if on:
|
||||
self._ensure_task()
|
||||
|
||||
async def _run(self) -> None:
|
||||
logger.info(f"[MONITOR] {self.unit_id}: feed started")
|
||||
try:
|
||||
while self._has_demand():
|
||||
snap, mst = await self._poll_once()
|
||||
if snap is not None:
|
||||
payload = _snapshot_payload(snap, self.unit_id, mst)
|
||||
self._broadcast(payload)
|
||||
try:
|
||||
await alert_evaluator.evaluate(self.unit_id, snap)
|
||||
except Exception as e:
|
||||
logger.warning(f"[MONITOR] {self.unit_id}: alert eval failed: {e}")
|
||||
await asyncio.sleep(MONITOR_POLL_INTERVAL)
|
||||
finally:
|
||||
logger.info(f"[MONITOR] {self.unit_id}: feed stopped")
|
||||
|
||||
async def _poll_once(self):
|
||||
"""One DOD poll: read, persist, return (snapshot, measurement_start_iso)."""
|
||||
db = SessionLocal()
|
||||
try:
|
||||
cfg = db.query(NL43Config).filter_by(unit_id=self.unit_id).first()
|
||||
if not cfg or not cfg.tcp_enabled:
|
||||
return None, None
|
||||
client = NL43Client(
|
||||
cfg.host, cfg.tcp_port,
|
||||
ftp_username=cfg.ftp_username, ftp_password=cfg.ftp_password,
|
||||
ftp_port=cfg.ftp_port or 21,
|
||||
)
|
||||
snap = await client.request_dod()
|
||||
snap.unit_id = self.unit_id
|
||||
persist_snapshot(snap, db)
|
||||
db.commit()
|
||||
status = db.query(NL43Status).filter_by(unit_id=self.unit_id).first()
|
||||
mst = (status.measurement_start_time.isoformat()
|
||||
if status and status.measurement_start_time else None)
|
||||
return snap, mst
|
||||
except Exception as e:
|
||||
logger.warning(f"[MONITOR] {self.unit_id}: poll failed: {e}")
|
||||
return None, None
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def _broadcast(self, payload: dict) -> None:
|
||||
for q in list(self._subscribers):
|
||||
try:
|
||||
q.put_nowait(payload)
|
||||
except asyncio.QueueFull:
|
||||
# Slow consumer — drop this frame rather than stall the whole feed.
|
||||
pass
|
||||
|
||||
|
||||
class MonitorManager:
|
||||
"""Registry of per-device monitors (one per unit_id)."""
|
||||
|
||||
def __init__(self):
|
||||
self._monitors: Dict[str, DeviceMonitor] = {}
|
||||
self._lock = asyncio.Lock()
|
||||
|
||||
async def get(self, unit_id: str) -> DeviceMonitor:
|
||||
async with self._lock:
|
||||
m = self._monitors.get(unit_id)
|
||||
if m is None:
|
||||
m = DeviceMonitor(unit_id)
|
||||
self._monitors[unit_id] = m
|
||||
return m
|
||||
|
||||
def status(self) -> dict:
|
||||
return {
|
||||
uid: {
|
||||
"running": m.running,
|
||||
"subscribers": m.subscriber_count(),
|
||||
"keepalive": m._keepalive,
|
||||
}
|
||||
for uid, m in self._monitors.items()
|
||||
}
|
||||
|
||||
|
||||
# Module-level singleton
|
||||
monitor_manager = MonitorManager()
|
||||
+315
-1
@@ -11,7 +11,7 @@ import os
|
||||
import asyncio
|
||||
|
||||
from app.database import get_db
|
||||
from app.models import NL43Config, NL43Status
|
||||
from app.models import NL43Config, NL43Status, AlertRule, AlertEvent
|
||||
from app.services import NL43Client, persist_snapshot
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -121,6 +121,310 @@ async def flush_connection_pool():
|
||||
return {"status": "ok", "message": "All cached connections closed"}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/disconnect")
|
||||
async def disconnect_device(unit_id: str, db: Session = Depends(get_db)):
|
||||
"""Cleanly close SLMM's persistent TCP connection to a single device.
|
||||
|
||||
Gracefully closes (TCP FIN + wait_closed) the pooled connection for this
|
||||
device and removes it from the pool, freeing the NL43's single connection
|
||||
slot. Idempotent — a no-op if no connection is currently cached.
|
||||
|
||||
Note: this releases the *idle* pooled connection. It does not interrupt an
|
||||
in-progress DRD stream or an in-flight command (those have the socket
|
||||
checked out of the pool) — close the stream WebSocket to end a live stream.
|
||||
"""
|
||||
cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first()
|
||||
if not cfg:
|
||||
raise HTTPException(status_code=404, detail="NL43 config not found")
|
||||
|
||||
from app.services import _connection_pool
|
||||
|
||||
device_key = f"{cfg.host}:{cfg.tcp_port}"
|
||||
had_conn = device_key in _connection_pool.get_stats().get("connections", {})
|
||||
|
||||
await _connection_pool.discard(device_key)
|
||||
|
||||
return {
|
||||
"status": "ok",
|
||||
"unit_id": unit_id,
|
||||
"device_key": device_key,
|
||||
"disconnected": had_conn,
|
||||
"message": "Connection closed" if had_conn else "No cached connection to close",
|
||||
}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/deactivate")
|
||||
async def deactivate_device(unit_id: str, db: Session = Depends(get_db)):
|
||||
"""Make a single unit dormant: stop background polling for it AND drop its
|
||||
connection, freeing the device's connection slot. poll_enabled=False is
|
||||
persisted, so the unit stays dormant across restarts until /activate.
|
||||
"""
|
||||
cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first()
|
||||
if not cfg:
|
||||
raise HTTPException(status_code=404, detail="NL43 config not found")
|
||||
|
||||
cfg.poll_enabled = False
|
||||
db.commit()
|
||||
|
||||
from app.services import _connection_pool, _get_device_lock
|
||||
|
||||
device_key = f"{cfg.host}:{cfg.tcp_port}"
|
||||
|
||||
# Wait briefly for any in-flight poll/command to finish (so its connection is
|
||||
# back in the pool), then drop it. If a long-lived stream holds the lock we
|
||||
# don't block forever — discard the pooled connection regardless.
|
||||
lock = await _get_device_lock(device_key)
|
||||
acquired = False
|
||||
try:
|
||||
await asyncio.wait_for(lock.acquire(), timeout=10.0)
|
||||
acquired = True
|
||||
except asyncio.TimeoutError:
|
||||
acquired = False
|
||||
try:
|
||||
await _connection_pool.discard(device_key)
|
||||
finally:
|
||||
if acquired:
|
||||
lock.release()
|
||||
|
||||
return {
|
||||
"status": "ok",
|
||||
"unit_id": unit_id,
|
||||
"poll_enabled": False,
|
||||
"message": "Polling disabled and connection closed for this unit",
|
||||
}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/activate")
|
||||
async def activate_device(unit_id: str, db: Session = Depends(get_db)):
|
||||
"""Resume background polling for a unit previously deactivated."""
|
||||
cfg = db.query(NL43Config).filter_by(unit_id=unit_id).first()
|
||||
if not cfg:
|
||||
raise HTTPException(status_code=404, detail="NL43 config not found")
|
||||
|
||||
cfg.poll_enabled = True
|
||||
db.commit()
|
||||
|
||||
return {
|
||||
"status": "ok",
|
||||
"unit_id": unit_id,
|
||||
"poll_enabled": True,
|
||||
"message": "Polling enabled for this unit",
|
||||
}
|
||||
|
||||
|
||||
@router.get("/_system/status")
|
||||
async def system_status():
|
||||
"""Report whether this SLMM instance is actively polling or in standby."""
|
||||
from app.background_poller import poller
|
||||
from app.services import _connection_pool
|
||||
return {
|
||||
"status": "ok",
|
||||
"mode": "active" if poller.is_active() else "standby",
|
||||
"polling_active": poller.is_active(),
|
||||
"active_connections": _connection_pool.get_stats().get("active_connections", 0),
|
||||
}
|
||||
|
||||
|
||||
@router.post("/_system/standby")
|
||||
async def system_standby():
|
||||
"""Put this SLMM instance into standby: stop polling ALL devices and release
|
||||
every connection, so it stops occupying device slots (e.g. so a prod instance
|
||||
can take over). Runtime-only — on restart the instance returns to its
|
||||
SLMM_POLLING_ENABLED default.
|
||||
"""
|
||||
from app.background_poller import poller
|
||||
await poller.set_active(False)
|
||||
return {"status": "ok", "mode": "standby",
|
||||
"message": "Polling stopped and all device connections released"}
|
||||
|
||||
|
||||
@router.post("/_system/resume")
|
||||
async def system_resume():
|
||||
"""Resume polling after standby (global)."""
|
||||
from app.background_poller import poller
|
||||
await poller.set_active(True)
|
||||
return {"status": "ok", "mode": "active", "message": "Polling resumed"}
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# LIVE MONITOR (fan-out) — one DOD feed per device, broadcast to many clients
|
||||
# ============================================================================
|
||||
|
||||
@router.websocket("/{unit_id}/monitor")
|
||||
async def monitor_stream(websocket: WebSocket, unit_id: str):
|
||||
"""Subscribe a browser to the device's shared 1 Hz DOD feed.
|
||||
|
||||
Any number of clients can attach without each opening its own device
|
||||
connection (one poll loop per device, fanned out). Same JSON shape as the
|
||||
DRD stream, but DOD-sourced so it includes ln1/ln2 (L1/L10).
|
||||
"""
|
||||
await websocket.accept()
|
||||
from app.monitor import monitor_manager
|
||||
|
||||
monitor = await monitor_manager.get(unit_id)
|
||||
queue = await monitor.subscribe()
|
||||
logger.info(f"Monitor subscriber attached for {unit_id} ({monitor.subscriber_count()} total)")
|
||||
|
||||
async def _watch_disconnect():
|
||||
# Completes when the client disconnects, so an idle feed (no data) still
|
||||
# detects the drop and we don't leak a subscription that keeps the device
|
||||
# feed (and its connection) alive.
|
||||
try:
|
||||
while True:
|
||||
msg = await websocket.receive()
|
||||
if msg.get("type") == "websocket.disconnect":
|
||||
return
|
||||
except Exception:
|
||||
return
|
||||
|
||||
gone = asyncio.ensure_future(_watch_disconnect())
|
||||
try:
|
||||
while not gone.done():
|
||||
try:
|
||||
payload = await asyncio.wait_for(queue.get(), timeout=1.0)
|
||||
except asyncio.TimeoutError:
|
||||
continue # re-check gone.done()
|
||||
await websocket.send_json(payload)
|
||||
except WebSocketDisconnect:
|
||||
logger.info(f"Monitor subscriber disconnected for {unit_id}")
|
||||
except Exception as e:
|
||||
logger.warning(f"Monitor stream error for {unit_id}: {e}")
|
||||
finally:
|
||||
gone.cancel()
|
||||
await monitor.unsubscribe(queue)
|
||||
|
||||
|
||||
@router.post("/{unit_id}/monitor/start")
|
||||
async def monitor_start(unit_id: str):
|
||||
"""Keep the device's feed running even with no browser attached, so alerting
|
||||
evaluates continuously. Runtime-only (resets on restart)."""
|
||||
from app.monitor import monitor_manager
|
||||
monitor = await monitor_manager.get(unit_id)
|
||||
await monitor.set_keepalive(True)
|
||||
return {"status": "ok", "unit_id": unit_id, "running": monitor.running, "keepalive": True}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/monitor/stop")
|
||||
async def monitor_stop(unit_id: str):
|
||||
"""Drop the keep-alive; the feed stops once no browser subscribers remain."""
|
||||
from app.monitor import monitor_manager
|
||||
monitor = await monitor_manager.get(unit_id)
|
||||
await monitor.set_keepalive(False)
|
||||
return {"status": "ok", "unit_id": unit_id, "keepalive": False}
|
||||
|
||||
|
||||
@router.get("/_monitor/status")
|
||||
async def monitor_status():
|
||||
"""Status of every device monitor (running, subscriber count, keep-alive)."""
|
||||
from app.monitor import monitor_manager
|
||||
return {"status": "ok", "monitors": monitor_manager.status()}
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# ALERTS — threshold rules + fired events
|
||||
# ============================================================================
|
||||
|
||||
class AlertRulePayload(BaseModel):
|
||||
name: str = "Alert"
|
||||
metric: str = "lp" # lp/leq/lmax/lmin/lpeak/ln1/ln2
|
||||
comparison: str = "above" # above | below
|
||||
threshold_db: float
|
||||
duration_s: int = 0 # sustained seconds before firing (0 = instant)
|
||||
clear_margin_db: float = 2.0 # hysteresis band
|
||||
cooldown_s: int = 300
|
||||
schedule_start: str | None = None # "HH:MM" local; null = always
|
||||
schedule_end: str | None = None
|
||||
schedule_days: str | None = None # CSV of 0-6 (Mon=0); null = every day
|
||||
channels: str = "log"
|
||||
recipients: str | None = None
|
||||
enabled: bool = True
|
||||
|
||||
|
||||
def _rule_dict(r: AlertRule) -> dict:
|
||||
return {
|
||||
"id": r.id, "unit_id": r.unit_id, "name": r.name, "metric": r.metric,
|
||||
"comparison": r.comparison, "threshold_db": r.threshold_db,
|
||||
"duration_s": r.duration_s, "clear_margin_db": r.clear_margin_db,
|
||||
"cooldown_s": r.cooldown_s, "schedule_start": r.schedule_start,
|
||||
"schedule_end": r.schedule_end, "schedule_days": r.schedule_days,
|
||||
"channels": r.channels, "recipients": r.recipients, "enabled": r.enabled,
|
||||
}
|
||||
|
||||
|
||||
def _event_dict(e: AlertEvent) -> dict:
|
||||
return {
|
||||
"id": e.id, "rule_id": e.rule_id, "unit_id": e.unit_id,
|
||||
"rule_name": e.rule_name, "metric": e.metric, "threshold_db": e.threshold_db,
|
||||
"onset_at": e.onset_at.isoformat() if e.onset_at else None,
|
||||
"onset_value": e.onset_value, "peak_value": e.peak_value,
|
||||
"clear_at": e.clear_at.isoformat() if e.clear_at else None,
|
||||
"status": e.status,
|
||||
"acknowledged_at": e.acknowledged_at.isoformat() if e.acknowledged_at else None,
|
||||
"acknowledged_by": e.acknowledged_by,
|
||||
}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/alerts/rules")
|
||||
def create_alert_rule(unit_id: str, payload: AlertRulePayload, db: Session = Depends(get_db)):
|
||||
rule = AlertRule(unit_id=unit_id, **payload.model_dump())
|
||||
db.add(rule)
|
||||
db.commit()
|
||||
db.refresh(rule)
|
||||
from app.alerts import alert_evaluator
|
||||
alert_evaluator.invalidate(unit_id)
|
||||
return {"status": "ok", "rule": _rule_dict(rule)}
|
||||
|
||||
|
||||
@router.get("/{unit_id}/alerts/rules")
|
||||
def list_alert_rules(unit_id: str, db: Session = Depends(get_db)):
|
||||
rules = db.query(AlertRule).filter_by(unit_id=unit_id).all()
|
||||
return {"status": "ok", "rules": [_rule_dict(r) for r in rules]}
|
||||
|
||||
|
||||
@router.put("/{unit_id}/alerts/rules/{rule_id}")
|
||||
def update_alert_rule(unit_id: str, rule_id: int, payload: AlertRulePayload, db: Session = Depends(get_db)):
|
||||
rule = db.query(AlertRule).filter_by(id=rule_id, unit_id=unit_id).first()
|
||||
if not rule:
|
||||
raise HTTPException(status_code=404, detail="Alert rule not found")
|
||||
for field, value in payload.model_dump().items():
|
||||
setattr(rule, field, value)
|
||||
db.commit()
|
||||
db.refresh(rule)
|
||||
from app.alerts import alert_evaluator
|
||||
alert_evaluator.invalidate(unit_id)
|
||||
return {"status": "ok", "rule": _rule_dict(rule)}
|
||||
|
||||
|
||||
@router.delete("/{unit_id}/alerts/rules/{rule_id}")
|
||||
def delete_alert_rule(unit_id: str, rule_id: int, db: Session = Depends(get_db)):
|
||||
rule = db.query(AlertRule).filter_by(id=rule_id, unit_id=unit_id).first()
|
||||
if not rule:
|
||||
raise HTTPException(status_code=404, detail="Alert rule not found")
|
||||
db.delete(rule)
|
||||
db.commit()
|
||||
from app.alerts import alert_evaluator
|
||||
alert_evaluator.invalidate(unit_id)
|
||||
return {"status": "ok", "deleted": rule_id}
|
||||
|
||||
|
||||
@router.get("/{unit_id}/alerts/events")
|
||||
def list_alert_events(unit_id: str, limit: int = 50, db: Session = Depends(get_db)):
|
||||
events = (db.query(AlertEvent).filter_by(unit_id=unit_id)
|
||||
.order_by(AlertEvent.onset_at.desc()).limit(limit).all())
|
||||
return {"status": "ok", "events": [_event_dict(e) for e in events]}
|
||||
|
||||
|
||||
@router.post("/{unit_id}/alerts/events/{event_id}/ack")
|
||||
def ack_alert_event(unit_id: str, event_id: int, by: str | None = None, db: Session = Depends(get_db)):
|
||||
evt = db.query(AlertEvent).filter_by(id=event_id, unit_id=unit_id).first()
|
||||
if not evt:
|
||||
raise HTTPException(status_code=404, detail="Alert event not found")
|
||||
evt.acknowledged_at = datetime.utcnow()
|
||||
evt.acknowledged_by = by
|
||||
db.commit()
|
||||
return {"status": "ok", "acknowledged": event_id}
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# GLOBAL POLLING STATUS ENDPOINT (must be before /{unit_id} routes)
|
||||
# ============================================================================
|
||||
@@ -450,6 +754,8 @@ def get_status(unit_id: str, db: Session = Depends(get_db)):
|
||||
"lmax": status.lmax,
|
||||
"lmin": status.lmin,
|
||||
"lpeak": status.lpeak,
|
||||
"ln1": status.ln1,
|
||||
"ln2": status.ln2,
|
||||
"battery_level": status.battery_level,
|
||||
"power_source": status.power_source,
|
||||
"sd_remaining_mb": status.sd_remaining_mb,
|
||||
@@ -472,6 +778,8 @@ class StatusPayload(BaseModel):
|
||||
lmax: str | None = None
|
||||
lmin: str | None = None
|
||||
lpeak: str | None = None
|
||||
ln1: str | None = None
|
||||
ln2: str | None = None
|
||||
battery_level: str | None = None
|
||||
power_source: str | None = None
|
||||
sd_remaining_mb: str | None = None
|
||||
@@ -504,6 +812,8 @@ def upsert_status(unit_id: str, payload: StatusPayload, db: Session = Depends(ge
|
||||
"lmax": status.lmax,
|
||||
"lmin": status.lmin,
|
||||
"lpeak": status.lpeak,
|
||||
"ln1": status.ln1,
|
||||
"ln2": status.ln2,
|
||||
"battery_level": status.battery_level,
|
||||
"power_source": status.power_source,
|
||||
"sd_remaining_mb": status.sd_remaining_mb,
|
||||
@@ -1205,6 +1515,8 @@ async def stream_live(websocket: WebSocket, unit_id: str):
|
||||
"lmax": snap.lmax, # Maximum level
|
||||
"lmin": snap.lmin, # Minimum level
|
||||
"lpeak": snap.lpeak, # Peak level
|
||||
"ln1": snap.ln1, # LN1 percentile (L1/L10 contract); null on DRD stream
|
||||
"ln2": snap.ln2, # LN2 percentile; null on DRD stream
|
||||
"raw_payload": snap.raw_payload,
|
||||
})
|
||||
except Exception as e:
|
||||
@@ -1876,6 +2188,8 @@ async def run_diagnostics(unit_id: str, db: Session = Depends(get_db)):
|
||||
"lmax": status.lmax,
|
||||
"lmin": status.lmin,
|
||||
"lpeak": status.lpeak,
|
||||
"ln1": status.ln1,
|
||||
"ln2": status.ln2,
|
||||
"battery_level": status.battery_level,
|
||||
"power_source": status.power_source,
|
||||
"sd_remaining_mb": status.sd_remaining_mb,
|
||||
|
||||
+68
-28
@@ -46,6 +46,8 @@ class NL43Snapshot:
|
||||
lmax: Optional[str] = None # Maximum level
|
||||
lmin: Optional[str] = None # Minimum level
|
||||
lpeak: Optional[str] = None # Peak level
|
||||
ln1: Optional[str] = None # Percentile slot LN1 (configurable; device default L5, contract L1)
|
||||
ln2: Optional[str] = None # Percentile slot LN2 (configurable; device default L10)
|
||||
battery_level: Optional[str] = None
|
||||
power_source: Optional[str] = None
|
||||
sd_remaining_mb: Optional[str] = None
|
||||
@@ -69,10 +71,27 @@ def persist_snapshot(s: NL43Snapshot, db: Session):
|
||||
|
||||
logger.info(f"State transition check for {s.unit_id}: '{previous_state}' -> '{new_state}'")
|
||||
|
||||
# Device returns "Start" when measuring, "Stop" when stopped
|
||||
# Normalize to previous behavior for backward compatibility
|
||||
is_measuring = new_state == "Start"
|
||||
was_measuring = previous_state == "Start"
|
||||
# The device reports "Start" while measuring; the DOD path uses that string,
|
||||
# but the DRD stream path tags snapshots "Measure" (and the DOD fallback also
|
||||
# uses "Measure"). Treat ALL of these as "measuring" — otherwise opening and
|
||||
# closing the live stream flips state "Start"->"Measure"->"Start", which the
|
||||
# old equality check misread as stop-then-start and reset measurement_start_time.
|
||||
#
|
||||
# Also: only act on RECOGNIZED states. A buffer desync on the shared connection
|
||||
# (e.g. right after a DRD/DOD test) can make a Measure? read return a stray,
|
||||
# garbled value; treating that as "not measuring" produced constant phantom
|
||||
# "STOPPED -> STARTED" log pairs and reset the timer. Ignore unknown reads.
|
||||
MEASURING_STATES = {"Start", "Measure"}
|
||||
STOPPED_STATES = {"Stop"}
|
||||
was_measuring = previous_state in MEASURING_STATES
|
||||
|
||||
if new_state in MEASURING_STATES:
|
||||
is_measuring = True
|
||||
elif new_state in STOPPED_STATES:
|
||||
is_measuring = False
|
||||
else:
|
||||
logger.warning(f"Ignoring unrecognized measurement state for {s.unit_id}: {new_state!r}")
|
||||
is_measuring = was_measuring # garbled/unknown read — no transition
|
||||
|
||||
if not was_measuring and is_measuring:
|
||||
# Measurement just started - record the start time
|
||||
@@ -95,13 +114,18 @@ def persist_snapshot(s: NL43Snapshot, db: Session):
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
row.measurement_state = new_state
|
||||
# Only persist a recognized state so one garbled read can't poison the next
|
||||
# transition check (which would manufacture the phantom STOPPED/STARTED pair).
|
||||
if new_state in MEASURING_STATES or new_state in STOPPED_STATES:
|
||||
row.measurement_state = new_state
|
||||
row.counter = s.counter
|
||||
row.lp = s.lp
|
||||
row.leq = s.leq
|
||||
row.lmax = s.lmax
|
||||
row.lmin = s.lmin
|
||||
row.lpeak = s.lpeak
|
||||
row.ln1 = s.ln1
|
||||
row.ln2 = s.ln2
|
||||
row.battery_level = s.battery_level
|
||||
row.power_source = s.power_source
|
||||
row.sd_remaining_mb = s.sd_remaining_mb
|
||||
@@ -338,14 +362,14 @@ class ConnectionPool:
|
||||
if self._is_alive(conn):
|
||||
self._drain_buffer(conn.reader)
|
||||
conn.last_used_at = time.time()
|
||||
logger.debug(f"Pool hit for {device_key} (age={time.time() - conn.created_at:.0f}s)")
|
||||
logger.info(f"Pool hit for {device_key} (age={time.time() - conn.created_at:.0f}s)")
|
||||
return conn.reader, conn.writer, True
|
||||
else:
|
||||
await self._close_connection(conn, reason="stale")
|
||||
|
||||
# Open fresh connection
|
||||
reader, writer = await self._open_connection(host, port, timeout)
|
||||
logger.debug(f"New connection opened for {device_key}")
|
||||
logger.info(f"New connection opened for {device_key}")
|
||||
return reader, writer, False
|
||||
|
||||
async def release(self, device_key: str, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, host: str, port: int):
|
||||
@@ -454,11 +478,11 @@ class ConnectionPool:
|
||||
"""Check whether a cached connection is still usable."""
|
||||
now = time.time()
|
||||
|
||||
# Age / idle checks
|
||||
if now - conn.last_used_at > self._idle_ttl:
|
||||
# Age / idle checks (value of -1 disables the check)
|
||||
if self._idle_ttl >= 0 and now - conn.last_used_at > self._idle_ttl:
|
||||
logger.debug(f"Connection {conn.device_key} idle too long ({now - conn.last_used_at:.0f}s > {self._idle_ttl}s)")
|
||||
return False
|
||||
if now - conn.created_at > self._max_age:
|
||||
if self._max_age >= 0 and now - conn.created_at > self._max_age:
|
||||
logger.debug(f"Connection {conn.device_key} too old ({now - conn.created_at:.0f}s > {self._max_age}s)")
|
||||
return False
|
||||
|
||||
@@ -691,22 +715,29 @@ class NL43Client:
|
||||
|
||||
snap = NL43Snapshot(unit_id="", raw_payload=resp, measurement_state=measurement_state)
|
||||
|
||||
# Parse known positions (based on NL43 communication guide - DRD format)
|
||||
# DRD format: d0=counter, d1=Lp, d2=Leq, d3=Lmax, d4=Lmin, d5=Lpeak, d6=LIeq, ...
|
||||
# Parse DOD positional fields. DOD's layout is DIFFERENT from DRD: it has NO
|
||||
# leading counter and it includes LE plus LN1–LN5. The device returns 4 channels
|
||||
# of 16 fields each — [Lp, Leq, LE, Lmax, Lmin, LN1, LN2, LN3, LN4, LN5, Lpeak,
|
||||
# LIeq, Leq_mov, Ltm5, over, under] — and channel 1 (parts[0:16]) is the main
|
||||
# display. The previous code reused the DRD map (treating parts[0] as a counter),
|
||||
# which shifted everything: Lp was reported as the counter, Leq as Lp, LE as Leq,
|
||||
# and LN1 as Lpeak (you could spot it because "Lpeak" came out < Lmax).
|
||||
try:
|
||||
# Capture d0 (counter) for timer synchronization
|
||||
if len(parts) >= 1:
|
||||
snap.counter = parts[0] # d0: Measurement interval counter (1-600)
|
||||
snap.lp = parts[0] # Lp: instantaneous sound pressure level
|
||||
if len(parts) >= 2:
|
||||
snap.lp = parts[1] # d1: Instantaneous sound pressure level
|
||||
if len(parts) >= 3:
|
||||
snap.leq = parts[2] # d2: Equivalent continuous sound level
|
||||
snap.leq = parts[1] # Leq: equivalent continuous level
|
||||
# parts[2] = LE (sound exposure level) — not currently surfaced
|
||||
if len(parts) >= 4:
|
||||
snap.lmax = parts[3] # d3: Maximum level
|
||||
snap.lmax = parts[3] # Lmax
|
||||
if len(parts) >= 5:
|
||||
snap.lmin = parts[4] # d4: Minimum level
|
||||
snap.lmin = parts[4] # Lmin
|
||||
if len(parts) >= 11:
|
||||
snap.lpeak = parts[10] # Lpeak (parts[5] is LN1, NOT Lpeak)
|
||||
if len(parts) >= 6:
|
||||
snap.lpeak = parts[5] # d5: Peak level
|
||||
snap.ln1 = parts[5] # LN1 percentile slot (device default L5; contract L1)
|
||||
if len(parts) >= 7:
|
||||
snap.ln2 = parts[6] # LN2 percentile slot (device default L10)
|
||||
except (IndexError, ValueError) as e:
|
||||
logger.warning(f"Error parsing DOD data points: {e}")
|
||||
|
||||
@@ -896,15 +927,20 @@ class NL43Client:
|
||||
# Acquire per-device lock - held for entire streaming session
|
||||
device_lock = await _get_device_lock(self.device_key)
|
||||
async with device_lock:
|
||||
# Evict any cached connection — streaming needs its own dedicated socket
|
||||
await _connection_pool.discard(self.device_key)
|
||||
await self._enforce_rate_limit()
|
||||
|
||||
logger.info(f"Starting DRD stream for {self.device_key}")
|
||||
|
||||
# Reuse the pooled connection instead of discard()+reopen. The NL43
|
||||
# allows only ONE TCP connection at a time, and on a cellular link the
|
||||
# device does not free its single slot fast enough for an immediate
|
||||
# reconnect — so a fresh connect times out (the DRD stream failure).
|
||||
# The per-device lock is held for the whole session, so it already
|
||||
# blocks the poller; reusing the warm socket keeps us at exactly one
|
||||
# connection and lets the stream start on the slot commands already use.
|
||||
try:
|
||||
reader, writer = await _connection_pool._open_connection(
|
||||
self.host, self.port, self.timeout
|
||||
reader, writer, from_cache = await _connection_pool.acquire(
|
||||
self.device_key, self.host, self.port, self.timeout
|
||||
)
|
||||
except ConnectionError:
|
||||
logger.error(f"DRD stream connection failed to {self.device_key}")
|
||||
@@ -981,16 +1017,20 @@ class NL43Client:
|
||||
break
|
||||
|
||||
finally:
|
||||
# Send SUB character to stop streaming
|
||||
# Stop streaming on the device (SUB = 0x1A), then return the warm
|
||||
# connection to the pool so subsequent commands reuse this single
|
||||
# socket instead of opening a second one. release() returns healthy
|
||||
# sockets to the pool and closes dead ones; the next acquire()
|
||||
# drains any residual stop output before reuse.
|
||||
try:
|
||||
writer.write(b"\x1A")
|
||||
await writer.drain()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
writer.close()
|
||||
with contextlib.suppress(Exception):
|
||||
await writer.wait_closed()
|
||||
await _connection_pool.release(
|
||||
self.device_key, reader, writer, self.host, self.port
|
||||
)
|
||||
|
||||
logger.info(f"DRD stream ended for {self.device_key}")
|
||||
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Migration script to add ln1 and ln2 percentile columns to the nl43_status table.
|
||||
|
||||
The NL-43 DOD response carries percentile slots LN1-LN5; the live SLM display
|
||||
(Terra-View) shows two of them (default L1/L10). This adds storage for the two
|
||||
surfaced slots. Run once per database to update existing schema.
|
||||
"""
|
||||
|
||||
import sqlite3
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
DB_PATH = Path(__file__).parent / "data" / "slmm.db"
|
||||
|
||||
|
||||
def migrate():
|
||||
"""Add ln1 and ln2 columns to the nl43_status table."""
|
||||
|
||||
if not DB_PATH.exists():
|
||||
print(f"Database not found at {DB_PATH}")
|
||||
print("No migration needed - database will be created with new schema")
|
||||
return
|
||||
|
||||
conn = sqlite3.connect(DB_PATH)
|
||||
cursor = conn.cursor()
|
||||
|
||||
try:
|
||||
cursor.execute("PRAGMA table_info(nl43_status)")
|
||||
columns = [row[1] for row in cursor.fetchall()]
|
||||
|
||||
if "ln1" in columns and "ln2" in columns:
|
||||
print("✓ ln1/ln2 columns already exist, no migration needed")
|
||||
return
|
||||
|
||||
if "ln1" not in columns:
|
||||
print("Adding ln1 column...")
|
||||
cursor.execute("ALTER TABLE nl43_status ADD COLUMN ln1 TEXT")
|
||||
print("✓ Added ln1 column")
|
||||
|
||||
if "ln2" not in columns:
|
||||
print("Adding ln2 column...")
|
||||
cursor.execute("ALTER TABLE nl43_status ADD COLUMN ln2 TEXT")
|
||||
print("✓ Added ln2 column")
|
||||
|
||||
conn.commit()
|
||||
print("\n✓ Migration completed successfully!")
|
||||
|
||||
except Exception as e:
|
||||
conn.rollback()
|
||||
print(f"✗ Migration failed: {e}", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
migrate()
|
||||
@@ -0,0 +1,68 @@
|
||||
"""
|
||||
Synthetic unit test for the alert state machine — no DB, no device.
|
||||
|
||||
Drives `_evaluate_step` with a fake clock + a level series and checks that
|
||||
onset/clear fire with the right debounce + hysteresis. Run:
|
||||
|
||||
docker compose exec -T slmm python3 test_alert_evaluator.py
|
||||
# or, if app.alerts imports cleanly standalone: python3 test_alert_evaluator.py
|
||||
"""
|
||||
|
||||
from types import SimpleNamespace
|
||||
from app.alerts import RuleState, _evaluate_step
|
||||
|
||||
|
||||
def rule(**kw):
|
||||
base = dict(threshold_db=85.0, duration_s=3, clear_margin_db=2.0, comparison="above")
|
||||
base.update(kw)
|
||||
return SimpleNamespace(**base)
|
||||
|
||||
|
||||
def run(series, r):
|
||||
st = RuleState()
|
||||
events = [(now, a) for value, now in series
|
||||
if (a := _evaluate_step(st, value, now, r))]
|
||||
return events, st
|
||||
|
||||
|
||||
def main():
|
||||
failures = 0
|
||||
|
||||
def check(label, cond, detail=""):
|
||||
nonlocal failures
|
||||
print(("PASS" if cond else "FAIL"), label, detail)
|
||||
if not cond:
|
||||
failures += 1
|
||||
|
||||
# 1) sustained exceedance -> onset after duration; recovery -> clear after duration
|
||||
r = rule(threshold_db=85, duration_s=3, clear_margin_db=2)
|
||||
ev, _ = run([(80, 0), (86, 1), (87, 2), (88, 3), (88, 4),
|
||||
(88, 5), (82, 6), (82, 7), (82, 8), (82, 9)], r)
|
||||
onsets = [t for t, a in ev if a == "onset"]
|
||||
clears = [t for t, a in ev if a == "clear"]
|
||||
check("1 sustained onset@4 / clear@9", onsets == [4] and clears == [9], str(ev))
|
||||
|
||||
# 2) brief spike under duration -> no onset (debounce)
|
||||
ev, _ = run([(80, 0), (90, 1), (90, 2), (80, 3), (80, 4)], rule(duration_s=3))
|
||||
check("2 brief spike debounced", ev == [], str(ev))
|
||||
|
||||
# 3) hysteresis: a dip into the margin (below threshold, above threshold-margin)
|
||||
# does NOT clear
|
||||
r = rule(threshold_db=85, duration_s=0, clear_margin_db=3)
|
||||
ev, st = run([(86, 0), (84, 1), (84, 2), (84, 3)], r)
|
||||
check("3 hysteresis holds ACTIVE", ev == [(0, "onset")] and st.phase == "active",
|
||||
f"{ev} phase={st.phase}")
|
||||
|
||||
# 4) 'below' comparison (device too quiet) -> onset when value < threshold
|
||||
ev, _ = run([(30, 0), (15, 1)], rule(threshold_db=20, duration_s=0,
|
||||
clear_margin_db=2, comparison="below"))
|
||||
check("4 below-comparison onset@1", ev == [(1, "onset")], str(ev))
|
||||
|
||||
print()
|
||||
print("ALL PASS" if failures == 0 else f"{failures} FAILURE(S)")
|
||||
return failures
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
import sys
|
||||
sys.exit(1 if main() else 0)
|
||||
Reference in New Issue
Block a user