Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8c17af4849 | |||
| b954eb8c89 | |||
| 0793e7df01 | |||
| 51dd6b682d | |||
| a7983d2958 | |||
| d6dd2e736b | |||
| af86cf713e | |||
| e3f9ca7f5b | |||
| ad1a40e0aa |
@@ -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:
|
||||
"""
|
||||
@@ -39,6 +45,7 @@ class BackgroundPoller:
|
||||
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."""
|
||||
@@ -71,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:
|
||||
|
||||
+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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
+135
@@ -121,6 +121,131 @@ 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"}
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# GLOBAL POLLING STATUS ENDPOINT (must be before /{unit_id} routes)
|
||||
# ============================================================================
|
||||
@@ -450,6 +575,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 +599,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 +633,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 +1336,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 +2009,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,
|
||||
|
||||
+63
-23
@@ -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
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user