From 0f829e03806c788c2a8c7661c85492fd2415ca27 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 15:12:34 +0000 Subject: [PATCH] Split into dedicated hardware poller + read-only consumers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rearchitect so a single poller process is the sole owner of all SAS/SMART/ SES hardware I/O, serialized behind one gate and paced, writing results to Redis. The API/web and MQTT publisher become pure Redis readers — they no longer issue any subprocess and can restart freely without touching the bus. This addresses backplane/expander stress from concurrent + restart-triggered SMART/SES storms (the prior model re-ran a hardware sweep on every container start, and polled sg_ses 0x02+0x07 every 60s; 0x07 errored on the IOM6/ Xyratex expanders). - poller.py: paced sweep (inventory, SMART per-drive w/ gap, SES 0x02, ZFS, host/MegaRAID), startup jitter, single-instance Redis lock, LED-queue worker - services/hwgate.py: global serialization semaphore (POLL_CONCURRENCY=1) - services/store.py: Redis as the only producer<->consumer interface + LED queue - services/health.py: shared drive-health classifier (fixes overview double-count) - gate smartctl/sg_ses/zpool/ledctl; drop sg_ses 0x07 from the hot path - routers + temps read-only from the store; main.py drops the in-process poller - compose: 3 services (privileged poller + unprivileged app + redis) w/ pacing knobs --- Dockerfile | 1 + README.md | 135 ++++++++++++++++-------------- docker-compose.yml | 52 ++++++++---- main.py | 97 ++++------------------ poller.py | 186 ++++++++++++++++++++++++++++++++++++++++++ routers/drives.py | 25 ++++-- routers/enclosures.py | 21 +++-- routers/leds.py | 21 +++-- routers/overview.py | 165 +++++++++++++------------------------ services/cache.py | 5 ++ services/enclosure.py | 38 ++++----- services/health.py | 21 +++++ services/host.py | 64 ++++++--------- services/hwgate.py | 27 ++++++ services/leds.py | 22 +++-- services/smart.py | 61 ++++++-------- services/store.py | 134 ++++++++++++++++++++++++++++++ services/temps.py | 81 ++++++++---------- services/zfs.py | 25 +++--- 19 files changed, 717 insertions(+), 464 deletions(-) create mode 100644 poller.py create mode 100644 services/health.py create mode 100644 services/hwgate.py create mode 100644 services/store.py diff --git a/Dockerfile b/Dockerfile index e702b6d..71cf3cc 100644 --- a/Dockerfile +++ b/Dockerfile @@ -24,6 +24,7 @@ COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY main.py . +COPY poller.py . COPY routers/ routers/ COPY services/ services/ COPY models/ models/ diff --git a/README.md b/README.md index 7162376..0733b4f 100644 --- a/README.md +++ b/README.md @@ -1,107 +1,118 @@ # JBOD Monitor -REST API for monitoring drive health in JBOD enclosures on Linux. +Drive-health monitoring for JBOD enclosures on Linux, with a REST API, web UI, +and optional Home Assistant MQTT publishing. -Auto-discovers SES enclosures via sysfs, maps drives to physical slots, and exposes SMART health data. +Auto-discovers SES enclosures via sysfs, maps drives to physical slots, reads +SMART, and aggregates per-enclosure temperatures (named SES sensors + hotspot). + +## Architecture + +Two processes, with **Redis as the only interface between them**: + +``` +poller (privileged, pid:host) ──smartctl/sg_ses/zpool/ledctl──> hardware + │ serialized + paced + ▼ + Redis ◀── reads ── app (REST API + web UI + MQTT) [unprivileged] + ▲ + └── locate-LED command queue (app enqueues, poller runs) +``` + +- **`poller`** is the *only* process that touches the SAS bus. All hardware + commands run behind a single gate (`POLL_CONCURRENCY`, default 1 = fully + serialized) with a gap between drives, so it never storms fragile SAS + expanders. It sweeps on an interval and writes results to Redis. +- **`app`** (and the MQTT publisher) are pure Redis readers — unprivileged, no + `/dev`. They can restart freely without issuing a single hardware command. +- **Locate LEDs**: the API enqueues a request in Redis; the poller executes it + serialized with everything else. + +This split exists for a reason: concurrent/bursty SMART+SES traffic — especially +repeated on every container restart — can drop SAS expanders and suspend pools. +One paced owner makes that structurally impossible. ## Prerequisites - Linux with SAS/SATA JBODs connected via HBA -- `smartmontools` — for `smartctl` (SMART data) -- `sg3-utils` — for `sg_ses` (SES enclosure data) +- `smartmontools` (`smartctl`), `sg3-utils` (`sg_ses`), `ledmon` (`ledctl`) +- Redis - Python 3.11+ ```bash -# Debian/Ubuntu -apt install smartmontools sg3-utils - -# RHEL/Fedora -dnf install smartmontools sg3_utils +apt install smartmontools sg3-utils ledmon # Debian/Ubuntu ``` -## Install +## Run (docker compose) ```bash -pip install -r requirements.txt +docker compose up -d --build ``` -## Run - -The API needs root access for `smartctl` to query drives: +Brings up `poller` (privileged), `app` (port 8000), and `redis`. To run the two +processes by hand instead: ```bash -sudo uvicorn main:app --host 0.0.0.0 --port 8000 +REDIS_HOST=localhost sudo -E python poller.py # producer (needs root) +REDIS_HOST=localhost uvicorn main:app --port 8000 # consumer ``` ## API Endpoints | Endpoint | Description | |---|---| -| `GET /api/health` | Service health + tool availability | -| `GET /api/enclosures` | List all discovered SES enclosures | -| `GET /api/enclosures/{id}/drives` | List drive slots for an enclosure | +| `GET /api/health` | Service health + `redis`/`mqtt`/`poller_fresh` status | +| `GET /api/enclosures` | List discovered SES enclosures | +| `GET /api/enclosures/{id}/drives` | Drive slots for an enclosure | | `GET /api/drives/{device}` | SMART detail for a block device | | `GET /api/overview` | Aggregate enclosure + drive health | | `GET /api/temps` | Per-enclosure temperatures: named SES sensors + hotspot | -| `GET /docs` | Interactive API docs (Swagger UI) | +| `POST /api/drives/{device}/led` | Queue a locate-LED change (`state`: `locate`/`off`) | +| `GET /docs` | Swagger UI | + +Read endpoints serve the poller's last sweep; `X-Poll-Age` (seconds) headers and +the `poller_fresh` health flag surface staleness. + +## Poller tuning (env) + +| Variable | Default | Description | +|---|---|---| +| `POLL_SWEEP_INTERVAL` | `300` | Seconds between full sweeps | +| `POLL_DRIVE_GAP` | `0.5` | Seconds between each drive's SMART query | +| `POLL_CONCURRENCY` | `1` | Max concurrent hardware ops (1 = serialized) | +| `POLL_STARTUP_JITTER` | `10` | Random startup delay to avoid restart storms | +| `POLL_DATA_TTL` | `900` | Redis TTL for poller data (must exceed sweep interval) | +| `ZFS_USE_NSENTER` | — | `true` to run `zpool` in the host namespace (containers) | ## Home Assistant (MQTT) -The monitor can publish per-enclosure temperatures to Home Assistant over MQTT -using [HA discovery](https://www.home-assistant.io/integrations/mqtt/#mqtt-discovery) — -no HA YAML required. Each enclosure becomes a device with: - -- a **Hotspot** sensor — the hottest reading across all SES temperature - sensors *and* the drives housed in that enclosure (drive temps are usually - the real hotspot); `drive_max_c` and `drive_temp_count` are exposed as - attributes. -- one sensor per named SES temperature element (e.g. *Temp Inlet*, - *Temp Outlet*), labelled from the SES element descriptor page. - -Publishing is **opt-in**: set `MQTT_HOST` to enable it. +The `app` publishes per-enclosure temperatures via +[MQTT discovery](https://www.home-assistant.io/integrations/mqtt/#mqtt-discovery) +— no HA YAML. Each enclosure becomes a device with a **Hotspot** sensor (max +across SES sensors + housed drive temps; `drive_max_c`/`drive_temp_count` as +attributes) and one sensor per named SES temperature element. Opt-in via +`MQTT_HOST`; availability is tracked via an MQTT LWT. | Variable | Default | Description | |---|---|---| | `MQTT_HOST` | _(unset)_ | Broker host. Unset → MQTT disabled. | | `MQTT_PORT` | `1883` | Broker port | -| `MQTT_USERNAME` / `MQTT_PASSWORD` | _(unset)_ | Broker credentials (fallback if OpenBao unused/unavailable) | +| `MQTT_USERNAME` / `MQTT_PASSWORD` | _(unset)_ | Broker credentials | | `MQTT_DISCOVERY_PREFIX` | `homeassistant` | HA discovery topic prefix | | `MQTT_BASE_TOPIC` | `jbod-monitor` | State/availability topic root | | `MQTT_PUBLISH_INTERVAL` | `60` | Seconds between publishes | -| `MQTT_NODE_ID` | _(hostname)_ | Stable id for this monitor (disambiguates multiple hosts) | -| `MQTT_CLIENT_ID` | `jbod-monitor-` | MQTT client id | +| `MQTT_NODE_ID` | _(hostname)_ | Stable id (disambiguates multiple hosts) | -### Credentials via OpenBao - -Broker credentials are sourced from OpenBao at startup rather than stored in -plaintext (matching the homelab runtime-secret pattern). Set `MQTT_SECRET_PATH` -to a KV v2 path holding `username`/`password`; the app reads -`/v1//data/` with the token from -`OPENBAO_TOKEN_FILE` (or `OPENBAO_TOKEN`). If the read fails, it falls back to -the `MQTT_USERNAME`/`MQTT_PASSWORD` env vars. - -| Variable | Default | Description | -|---|---|---| -| `MQTT_SECRET_PATH` | _(unset)_ | KV v2 path holding the broker creds (e.g. `home_assistant`). Unset → use env creds. | -| `MQTT_SECRET_USER_KEY` | `username` | Key within the secret for the broker username | -| `MQTT_SECRET_PASS_KEY` | `password` | Key within the secret for the broker password | -| `OPENBAO_ADDR` | `https://vault.adamksmith.xyz` | OpenBao API base URL | -| `OPENBAO_MOUNT` | `secret` | KV v2 mount point | -| `OPENBAO_TOKEN_FILE` | _(unset)_ | Path to a file containing the OpenBao token (preferred; mount read-only) | -| `OPENBAO_TOKEN` | _(unset)_ | OpenBao token (used if no token file) | - -The deployed compose reads the broker user/pass from `secret/home_assistant` -(keys `mqtt_user`/`mqtt_pass`) using the read-only `claude-read` token mounted -at `/run/secrets/openbao-token`. - -Topics (defaults): +**Credentials**: simplest is a root-owned `/etc/jbod-monitor/mqtt.env` with +`MQTT_USERNAME=`/`MQTT_PASSWORD=`, loaded via compose `env_file`. Optionally the +app can read them from OpenBao at runtime (`MQTT_SECRET_PATH` + `OPENBAO_*`, +keys via `MQTT_SECRET_USER_KEY`/`MQTT_SECRET_PASS_KEY`) — see +`services/secrets.py`; it falls back to the env vars if the read fails. +Topics: ``` homeassistant/sensor/_enc/hotspot/config # retained discovery homeassistant/sensor/_enc/temp/config # retained discovery jbod-monitor//status # online | offline (LWT) -jbod-monitor//enclosure//state # {"hotspot_c":42.0, "sensor_0":28.5, ...} +jbod-monitor//enclosure//state # {"hotspot_c":42.0, ...} ``` - -Availability is tracked via an MQTT LWT, so entities show *unavailable* in HA -if the monitor stops. diff --git a/docker-compose.yml b/docker-compose.yml index 6ab1226..4b1f81e 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,43 +1,59 @@ services: - jbod-monitor: + # Producer: the ONLY service that touches hardware (smartctl/sg_ses/zpool/ + # ledctl). Serialized + paced; writes everything to Redis. + poller: build: . image: docker.adamksmith.xyz/jbod-monitor:latest - container_name: jbod-monitor + container_name: jbod-poller restart: unless-stopped privileged: true pid: host network_mode: host + command: ["python", "-u", "poller.py"] volumes: - /dev:/dev - /sys:/sys:ro - /run/udev:/run/udev:ro - # OpenBao RO token (claude-read) mounted read-only; never bake into image. - - /etc/jbod-monitor/openbao-token:/run/secrets/openbao-token:ro environment: - TZ=America/Denver - - UVICORN_LOG_LEVEL=info - ZFS_USE_NSENTER=true - REDIS_HOST=127.0.0.1 - REDIS_PORT=6379 - REDIS_DB=0 - - SMART_CACHE_TTL=120 - - SMART_POLL_INTERVAL=90 - # Home Assistant MQTT publishing (opt-in: leave MQTT_HOST unset to disable). - # EMQX at 10.5.30.3 is reachable by IP and bridged to the Mosquitto HA - # add-on, so discovery reaches Home Assistant either way. + # Pacing — keep the SAS bus quiet. POLL_CONCURRENCY=1 fully serializes. + - POLL_SWEEP_INTERVAL=300 + - POLL_DRIVE_GAP=0.5 + - POLL_CONCURRENCY=1 + - POLL_STARTUP_JITTER=10 + - POLL_DATA_TTL=900 + depends_on: + - redis + + # Consumer: REST API + web UI + Home Assistant MQTT. Reads Redis only — + # unprivileged, no /dev, can restart freely without touching hardware. + app: + build: . + image: docker.adamksmith.xyz/jbod-monitor:latest + container_name: jbod-monitor + restart: unless-stopped + network_mode: host + environment: + - TZ=America/Denver + - UVICORN_LOG_LEVEL=info + - REDIS_HOST=127.0.0.1 + - REDIS_PORT=6379 + - REDIS_DB=0 + # Home Assistant MQTT publishing (EMQX, bridged to Mosquitto HA add-on). - MQTT_HOST=10.5.30.3 - MQTT_PORT=1883 - MQTT_DISCOVERY_PREFIX=homeassistant - MQTT_BASE_TOPIC=jbod-monitor - MQTT_PUBLISH_INTERVAL=60 - # Broker credentials sourced from OpenBao (secret/home_assistant, - # keys: mqtt_user/mqtt_pass). Falls back to MQTT_USERNAME/MQTT_PASSWORD - # env if the read fails. Token is the read-only claude-read token. - - OPENBAO_ADDR=https://vault.adamksmith.xyz - - OPENBAO_TOKEN_FILE=/run/secrets/openbao-token - - MQTT_SECRET_PATH=home_assistant - - MQTT_SECRET_USER_KEY=mqtt_user - - MQTT_SECRET_PASS_KEY=mqtt_pass + # Broker credentials (MQTT_USERNAME/MQTT_PASSWORD). Optional: container + # starts without it (MQTT just won't authenticate). + env_file: + - path: /etc/jbod-monitor/mqtt.env + required: false depends_on: - redis diff --git a/main.py b/main.py index bcb0756..f1ac65b 100644 --- a/main.py +++ b/main.py @@ -1,5 +1,4 @@ import asyncio -import json import logging import os from contextlib import asynccontextmanager @@ -11,11 +10,9 @@ from fastapi.staticfiles import StaticFiles from models.schemas import HealthCheck from routers import drives, enclosures, leds, overview, temps -from services.cache import cache_set, close_cache, init_cache, redis_available -from services.enclosure import discover_enclosures, list_slots +from services import store +from services.cache import close_cache, init_cache, redis_available from services.mqtt_publisher import MqttPublisher, _PAHO_AVAILABLE, mqtt_enabled -from services.smart import SMART_CACHE_TTL, _run_smartctl, sg_ses_available, smartctl_available -from services.zfs import get_zfs_pool_map logging.basicConfig( level=logging.INFO, @@ -23,81 +20,20 @@ logging.basicConfig( ) logger = logging.getLogger(__name__) -SMART_POLL_INTERVAL = int(os.environ.get("SMART_POLL_INTERVAL", "90")) - _tool_status: dict[str, bool] = {} -_poll_task: asyncio.Task | None = None _mqtt_task: asyncio.Task | None = None _mqtt_publisher: MqttPublisher | None = None -async def smart_poll_loop(): - """Pre-warm Redis with SMART data for all drives.""" - await asyncio.sleep(2) # let app finish starting - while True: - try: - # Discover all enclosure devices - enclosures_raw = discover_enclosures() - devices: set[str] = set() - for enc in enclosures_raw: - for slot in list_slots(enc["id"]): - if slot["device"]: - devices.add(slot["device"]) - - # Discover host block devices via lsblk - try: - proc = await asyncio.create_subprocess_exec( - "lsblk", "-d", "-o", "NAME,TYPE", "-J", - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, _ = await proc.communicate() - if stdout: - for dev in json.loads(stdout).get("blockdevices", []): - if dev.get("type") == "disk": - name = dev.get("name", "") - if name and name not in devices: - devices.add(name) - except Exception as e: - logger.warning("lsblk discovery failed in poller: %s", e) - - # Poll each drive and cache result - for device in sorted(devices): - try: - result = await _run_smartctl(device) - await cache_set(f"jbod:smart:{device}", result, SMART_CACHE_TTL) - except Exception as e: - logger.warning("Poll failed for %s: %s", device, e) - - # Pre-warm ZFS map (bypasses cache by calling directly) - await get_zfs_pool_map() - - logger.info("SMART poll complete: %d devices", len(devices)) - except Exception as e: - logger.error("SMART poll loop error: %s", e) - await asyncio.sleep(SMART_POLL_INTERVAL) - - @asynccontextmanager async def lifespan(app: FastAPI): - global _poll_task, _mqtt_task, _mqtt_publisher - # Startup - _tool_status["smartctl"] = smartctl_available() - _tool_status["sg_ses"] = sg_ses_available() - - if not _tool_status["smartctl"]: - logger.warning("smartctl not found — install smartmontools for SMART data") - if not _tool_status["sg_ses"]: - logger.warning("sg_ses not found — install sg3-utils for enclosure SES data") - if os.geteuid() != 0: - logger.warning("Not running as root — smartctl may fail on some devices") + """API/web + MQTT consumer. Reads everything from Redis; the dedicated + poller process is the only thing that touches hardware.""" + global _mqtt_task, _mqtt_publisher await init_cache() _tool_status["redis"] = redis_available() - if redis_available(): - _poll_task = asyncio.create_task(smart_poll_loop()) - # Optional MQTT publisher for Home Assistant (opt-in via MQTT_HOST). _tool_status["mqtt"] = False if mqtt_enabled(): @@ -114,7 +50,6 @@ async def lifespan(app: FastAPI): yield - # Shutdown if _mqtt_task is not None: _mqtt_task.cancel() try: @@ -123,19 +58,13 @@ async def lifespan(app: FastAPI): pass if _mqtt_publisher is not None: _mqtt_publisher.stop() - if _poll_task is not None: - _poll_task.cancel() - try: - await _poll_task - except asyncio.CancelledError: - pass await close_cache() app = FastAPI( title="JBOD Monitor", description="Drive health monitoring for JBOD enclosures", - version="0.1.0", + version="0.2.0", lifespan=lifespan, ) @@ -155,8 +84,18 @@ app.include_router(temps.router) @app.get("/api/health", response_model=HealthCheck, tags=["health"]) async def health(): - _tool_status["redis"] = redis_available() - return HealthCheck(status="ok", tools=_tool_status) + tools = dict(_tool_status) + tools["redis"] = redis_available() + + # Poller freshness: is the producer keeping the store warm? + meta = await store.get_meta() + if meta and meta.get("last_sweep_ts"): + age = store.now() - meta["last_sweep_ts"] + tools["poller_fresh"] = age < meta.get("sweep_interval", 300) * 3 + else: + tools["poller_fresh"] = False + + return HealthCheck(status="ok", tools=tools) # Serve built frontend static files — mounted last so /api routes take priority. diff --git a/poller.py b/poller.py new file mode 100644 index 0000000..5bff8a5 --- /dev/null +++ b/poller.py @@ -0,0 +1,186 @@ +"""Dedicated hardware poller — the sole owner of SAS/SMART/SES I/O. + +Runs as its own (privileged, pid:host) container. Everything that touches +the bus happens here, serialized behind the hardware gate and paced, then +written to Redis. The API and MQTT services are pure Redis readers and +never issue a hardware command — so they can crash-loop freely without +ever storming the backplanes. + +Design points that matter for fragile SAS expanders: + - one global semaphore (POLL_CONCURRENCY=1) serializes all hardware I/O + - a gap (POLL_DRIVE_GAP) between each drive keeps the bus mostly quiet + - startup jitter avoids a thundering herd on (re)start + - a single-instance Redis lock means two pollers can't both poll + - last-good results stay in Redis; a failed sweep never blanks them +""" +import asyncio +import logging +import os +import random +import socket +import time + +from services import store +from services.cache import close_cache, init_cache, redis_available +from services.enclosure import discover_enclosures, get_enclosure_status, list_slots +from services.host import get_host_drives +from services.leds import run_led +from services.smart import _run_smartctl, sg_ses_available, smartctl_available +from services.zfs import get_zfs_pool_map + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", +) +logger = logging.getLogger("poller") + +SWEEP_INTERVAL = int(os.environ.get("POLL_SWEEP_INTERVAL", "300")) +DRIVE_GAP = float(os.environ.get("POLL_DRIVE_GAP", "0.5")) +STARTUP_JITTER = float(os.environ.get("POLL_STARTUP_JITTER", "10")) +LOCK_TTL = 30 + +TOKEN = f"{socket.gethostname()}-{os.getpid()}" + + +async def sweep() -> None: + """One full, paced pass over all hardware → Redis.""" + t0 = time.monotonic() + errors: list[str] = [] + + # 1. Discover enclosures + slots (sysfs, no bus I/O). + enclosures = discover_enclosures() + inv_enclosures = [] + drive_devices: list[str] = [] + for enc in enclosures: + slots = list_slots(enc["id"]) + inv_enclosures.append({**enc, "slots": slots}) + for s in slots: + if s["device"]: + drive_devices.append(s["device"]) + + # 2. ZFS topology (gated). + pool_map = await store.get_zfs_map() + try: + pool_map = await get_zfs_pool_map() + await store.set_zfs_map(pool_map) + except Exception as e: + errors.append(f"zfs:{e}") + + # 3. Per-drive SMART, one at a time with a gap between each. + for dev in drive_devices: + try: + data = await _run_smartctl(dev) + await store.set_smart(dev, data) + except Exception as e: + errors.append(f"smart:{dev}:{e}") + await asyncio.sleep(DRIVE_GAP) + + # 4. Per-enclosure SES status (gated), paced. + for enc in enclosures: + if not enc.get("sg_device"): + continue + try: + ses = await get_enclosure_status(enc["sg_device"]) + if ses is not None: + await store.set_ses(enc["id"], ses) + except Exception as e: + errors.append(f"ses:{enc['id']}:{e}") + await asyncio.sleep(DRIVE_GAP) + + # 5. Host (non-enclosure) drives + MegaRAID (gated, reuses pool_map). + try: + host_drives = await get_host_drives(pool_map) + await store.set_host_drives(host_drives) + except Exception as e: + errors.append(f"host:{e}") + + # 6. Inventory + heartbeat. + await store.set_inventory({"enclosures": inv_enclosures, "ts": store.now()}) + duration = round(time.monotonic() - t0, 1) + await store.set_meta({ + "last_sweep_ts": store.now(), + "last_sweep_duration": duration, + "drive_count": len(drive_devices), + "enclosure_count": len(enclosures), + "errors": errors[:20], + "sweep_interval": SWEEP_INTERVAL, + }) + logger.info( + "sweep done: %d drives, %d enclosures, %.1fs, %d error(s)", + len(drive_devices), len(enclosures), duration, len(errors), + ) + + +async def led_consumer() -> None: + """Drain locate-LED requests from Redis and run them (gated).""" + while True: + try: + cmd = await store.pop_led(timeout=5) + if cmd: + try: + await run_led(cmd["device"], cmd["state"]) + logger.info("LED %s -> %s", cmd.get("device"), cmd.get("state")) + except Exception as e: + logger.warning("LED command failed %s: %s", cmd, e) + except asyncio.CancelledError: + raise + except Exception as e: + logger.warning("LED consumer error: %s", e) + await asyncio.sleep(2) + + +async def lock_refresher() -> None: + while True: + await asyncio.sleep(LOCK_TTL / 3) + if not await store.refresh_lock(TOKEN, LOCK_TTL): + logger.error("lost poller lock — exiting to avoid double-polling") + os._exit(1) + + +async def main() -> None: + await init_cache() + while not redis_available(): + logger.warning("Redis unavailable — poller needs Redis; retrying in 5s") + await asyncio.sleep(5) + await init_cache() + + # Single-instance guard. + while not await store.acquire_lock(TOKEN, LOCK_TTL): + logger.warning("another poller holds the lock; waiting %ds", LOCK_TTL // 2) + await asyncio.sleep(LOCK_TTL // 2) + logger.info("poller lock acquired (%s)", TOKEN) + + if not smartctl_available(): + logger.warning("smartctl not found — install smartmontools") + if not sg_ses_available(): + logger.warning("sg_ses not found — install sg3-utils") + if os.geteuid() != 0: + logger.warning("not running as root — smartctl/sg_ses may fail") + + # Stagger startup so a restart doesn't immediately storm the bus. + jitter = random.uniform(0, STARTUP_JITTER) + logger.info("startup jitter %.1fs", jitter) + await asyncio.sleep(jitter) + + refresher = asyncio.create_task(lock_refresher()) + leds = asyncio.create_task(led_consumer()) + try: + while True: + t0 = time.monotonic() + try: + await sweep() + except Exception as e: + logger.error("sweep error: %s", e) + elapsed = time.monotonic() - t0 + await asyncio.sleep(max(5, SWEEP_INTERVAL - elapsed)) + finally: + refresher.cancel() + leds.cancel() + await close_cache() + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + pass diff --git a/routers/drives.py b/routers/drives.py index ff38fc7..47ff0b2 100644 --- a/routers/drives.py +++ b/routers/drives.py @@ -1,30 +1,37 @@ +import re + from fastapi import APIRouter, HTTPException, Response from models.schemas import DriveDetail -from services.smart import get_smart_data -from services.zfs import get_zfs_pool_map +from services import store router = APIRouter(prefix="/api/drives", tags=["drives"]) @router.get("/{device}", response_model=DriveDetail) async def get_drive_detail(device: str, response: Response): - """Get SMART detail for a specific block device.""" - try: - data, cache_hit = await get_smart_data(device) - except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) + """Get SMART detail for a specific block device (from the store).""" + if not re.match(r"^[a-zA-Z0-9\-]+$", device): + raise HTTPException(status_code=400, detail=f"Invalid device name: {device}") + data = await store.get_smart(device) + if data is None: + raise HTTPException( + status_code=404, + detail=f"No SMART data for '{device}' yet (poller may not have swept it)", + ) if "error" in data: raise HTTPException(status_code=502, detail=data["error"]) - pool_map = await get_zfs_pool_map() + pool_map = await store.get_zfs_map() zfs_info = pool_map.get(device) if zfs_info: data["zfs_pool"] = zfs_info["pool"] data["zfs_vdev"] = zfs_info["vdev"] data["zfs_state"] = zfs_info.get("state") - response.headers["X-Cache"] = "HIT" if cache_hit else "MISS" + meta = await store.get_meta() + if meta and meta.get("last_sweep_ts"): + response.headers["X-Poll-Age"] = str(int(store.now() - meta["last_sweep_ts"])) return DriveDetail(**data) diff --git a/routers/enclosures.py b/routers/enclosures.py index 59858f5..2030163 100644 --- a/routers/enclosures.py +++ b/routers/enclosures.py @@ -1,24 +1,23 @@ from fastapi import APIRouter, HTTPException from models.schemas import Enclosure, SlotInfo -from services.enclosure import discover_enclosures, list_slots +from services import store router = APIRouter(prefix="/api/enclosures", tags=["enclosures"]) @router.get("", response_model=list[Enclosure]) async def get_enclosures(): - """Discover all SES enclosures.""" - return discover_enclosures() + """List all discovered SES enclosures (from the poller inventory).""" + inventory = await store.get_inventory() + return [Enclosure(**e) for e in inventory.get("enclosures", [])] @router.get("/{enclosure_id}/drives", response_model=list[SlotInfo]) async def get_enclosure_drives(enclosure_id: str): - """List all drive slots for an enclosure.""" - slots = list_slots(enclosure_id) - if not slots: - # Check if the enclosure exists at all - enclosures = discover_enclosures() - if not any(e["id"] == enclosure_id for e in enclosures): - raise HTTPException(status_code=404, detail=f"Enclosure '{enclosure_id}' not found") - return slots + """List all drive slots for an enclosure (from the poller inventory).""" + inventory = await store.get_inventory() + for enc in inventory.get("enclosures", []): + if enc["id"] == enclosure_id: + return [SlotInfo(**s) for s in enc.get("slots", [])] + raise HTTPException(status_code=404, detail=f"Enclosure '{enclosure_id}' not found") diff --git a/routers/leds.py b/routers/leds.py index 4ee8ec4..e7f0e29 100644 --- a/routers/leds.py +++ b/routers/leds.py @@ -1,7 +1,9 @@ +import re + from fastapi import APIRouter, HTTPException from pydantic import BaseModel, field_validator -from services.leds import set_led +from services import store router = APIRouter(prefix="/api/drives", tags=["leds"]) @@ -19,11 +21,12 @@ class LedRequest(BaseModel): @router.post("/{device}/led") async def set_drive_led(device: str, body: LedRequest): - """Set the locate LED for a drive.""" - try: - result = await set_led(device, body.state) - except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) - except RuntimeError as e: - raise HTTPException(status_code=502, detail=str(e)) - return result + """Queue a locate-LED change. The poller (sole hardware owner) executes it.""" + if not re.fullmatch(r"[a-zA-Z0-9]+", device): + raise HTTPException(status_code=400, detail=f"Invalid device name: {device}") + + queued = await store.enqueue_led(device, body.state) + if not queued: + raise HTTPException(status_code=503, detail="LED queue unavailable (Redis down)") + + return {"device": device, "state": body.state, "queued": True} diff --git a/routers/overview.py b/routers/overview.py index ca27c8e..3a00c8e 100644 --- a/routers/overview.py +++ b/routers/overview.py @@ -1,4 +1,3 @@ -import asyncio import logging from fastapi import APIRouter, Response @@ -11,10 +10,8 @@ from models.schemas import ( Overview, SlotWithDrive, ) -from services.enclosure import discover_enclosures, get_enclosure_status, list_slots -from services.host import get_host_drives -from services.smart import get_smart_data -from services.zfs import get_zfs_pool_map +from services import store +from services.health import compute_health_status logger = logging.getLogger(__name__) @@ -23,142 +20,92 @@ router = APIRouter(prefix="/api/overview", tags=["overview"]) @router.get("", response_model=Overview) async def get_overview(response: Response): - """Aggregate view of all enclosures, slots, and drive health.""" - enclosures_raw = discover_enclosures() - pool_map = await get_zfs_pool_map() + """Aggregate view of all enclosures, slots, and drive health. - # Fetch SES health data for all enclosures concurrently - async def _get_health(enc): - if enc.get("sg_device"): - return await get_enclosure_status(enc["sg_device"]) - return None - - health_results = await asyncio.gather( - *[_get_health(enc) for enc in enclosures_raw], - return_exceptions=True, - ) + Pure read from the store — the poller is the only thing that touched + hardware to produce this data. + """ + inventory = await store.get_inventory() + pool_map = await store.get_zfs_map() enc_results: list[EnclosureWithDrives] = [] total_drives = 0 warnings = 0 errors = 0 all_healthy = True - all_cache_hits = True - any_lookups = False - - for enc_idx, enc in enumerate(enclosures_raw): - slots_raw = list_slots(enc["id"]) - - # Gather SMART data for all populated slots concurrently - populated = [(s, s["device"]) for s in slots_raw if s["populated"] and s["device"]] - smart_tasks = [get_smart_data(dev) for _, dev in populated] - smart_results = await asyncio.gather(*smart_tasks, return_exceptions=True) - - smart_map: dict[str, dict] = {} - for (slot_info, dev), result in zip(populated, smart_results): - if isinstance(result, Exception): - logger.warning("SMART query failed for %s: %s", dev, result) - smart_map[dev] = {"device": dev, "smart_supported": False} - all_cache_hits = False - else: - data, hit = result - smart_map[dev] = data - any_lookups = True - if not hit: - all_cache_hits = False + for enc in inventory.get("enclosures", []): slots_out: list[SlotWithDrive] = [] - for s in slots_raw: + for s in enc.get("slots", []): + dev = s.get("device") drive_summary = None - if s["device"] and s["device"] in smart_map: - sd = smart_map[s["device"]] + + if dev: total_drives += 1 + sd = await store.get_smart(dev) + if sd: + status = compute_health_status(sd) + if status == "error": + errors += 1 + all_healthy = False + elif status == "warning": + warnings += 1 - healthy = sd.get("smart_healthy") - if healthy is False: - errors += 1 - all_healthy = False - elif healthy is None and sd.get("smart_supported", True): - warnings += 1 - - # Check for concerning SMART values - if sd.get("reallocated_sectors") and sd["reallocated_sectors"] > 0: - warnings += 1 - if sd.get("pending_sectors") and sd["pending_sectors"] > 0: - warnings += 1 - if sd.get("uncorrectable_errors") and sd["uncorrectable_errors"] > 0: - warnings += 1 - - # Compute health_status for frontend - realloc = sd.get("reallocated_sectors") or 0 - pending = sd.get("pending_sectors") or 0 - unc = sd.get("uncorrectable_errors") or 0 - if healthy is False: - health_status = "error" - elif realloc > 0 or pending > 0 or unc > 0 or (healthy is None and sd.get("smart_supported", True)): - health_status = "warning" - else: - health_status = "healthy" - - drive_summary = DriveHealthSummary( - device=sd["device"], - model=sd.get("model"), - serial=sd.get("serial"), - wwn=sd.get("wwn"), - firmware=sd.get("firmware"), - capacity_bytes=sd.get("capacity_bytes"), - smart_healthy=healthy, - smart_supported=sd.get("smart_supported", True), - temperature_c=sd.get("temperature_c"), - power_on_hours=sd.get("power_on_hours"), - reallocated_sectors=sd.get("reallocated_sectors"), - pending_sectors=sd.get("pending_sectors"), - uncorrectable_errors=sd.get("uncorrectable_errors"), - zfs_pool=pool_map.get(sd["device"], {}).get("pool"), - zfs_vdev=pool_map.get(sd["device"], {}).get("vdev"), - zfs_state=pool_map.get(sd["device"], {}).get("state"), - health_status=health_status, - ) - elif s["populated"]: + zinfo = pool_map.get(dev, {}) + drive_summary = DriveHealthSummary( + device=sd.get("device", dev), + model=sd.get("model"), + serial=sd.get("serial"), + wwn=sd.get("wwn"), + firmware=sd.get("firmware"), + capacity_bytes=sd.get("capacity_bytes"), + smart_healthy=sd.get("smart_healthy"), + smart_supported=sd.get("smart_supported", True), + temperature_c=sd.get("temperature_c"), + power_on_hours=sd.get("power_on_hours"), + reallocated_sectors=sd.get("reallocated_sectors"), + pending_sectors=sd.get("pending_sectors"), + uncorrectable_errors=sd.get("uncorrectable_errors"), + zfs_pool=zinfo.get("pool"), + zfs_vdev=zinfo.get("vdev"), + zfs_state=zinfo.get("state"), + health_status=status, + ) + elif s.get("populated"): total_drives += 1 slots_out.append(SlotWithDrive( slot=s["slot"], populated=s["populated"], - device=s["device"], + device=dev, drive=drive_summary, )) - # Attach enclosure health from SES - health_data = health_results[enc_idx] + ses = await store.get_ses(enc["id"]) enc_health = None - if isinstance(health_data, dict): - enc_health = EnclosureHealth(**health_data) - # Count enclosure-level issues + if ses: + enc_health = EnclosureHealth(**ses) if enc_health.overall_status == "CRITICAL": errors += 1 all_healthy = False elif enc_health.overall_status == "WARNING": warnings += 1 - elif isinstance(health_data, Exception): - logger.warning("SES health failed for %s: %s", enc["id"], health_data) enc_results.append(EnclosureWithDrives( id=enc["id"], sg_device=enc.get("sg_device"), - vendor=enc["vendor"], - model=enc["model"], - revision=enc["revision"], - total_slots=enc["total_slots"], - populated_slots=enc["populated_slots"], + vendor=enc.get("vendor", ""), + model=enc.get("model", ""), + revision=enc.get("revision", ""), + total_slots=enc.get("total_slots", 0), + populated_slots=enc.get("populated_slots", 0), slots=slots_out, health=enc_health, )) - # Host drives (non-enclosure) - host_drives_raw = await get_host_drives() + # Host (non-enclosure) drives — fully precomputed by the poller. host_drives_out: list[HostDrive] = [] - for hd in host_drives_raw: + for hd in await store.get_host_drives(): total_drives += 1 hs = hd.get("health_status", "healthy") if hs == "error": @@ -166,7 +113,6 @@ async def get_overview(response: Response): all_healthy = False elif hs == "warning": warnings += 1 - # Count physical drives behind RAID controllers for pd in hd.get("physical_drives", []): total_drives += 1 pd_hs = pd.get("health_status", "healthy") @@ -177,7 +123,10 @@ async def get_overview(response: Response): warnings += 1 host_drives_out.append(HostDrive(**hd)) - response.headers["X-Cache"] = "HIT" if (any_lookups and all_cache_hits) else "MISS" + # Surface poller freshness so stale data is visible. + meta = await store.get_meta() + if meta and meta.get("last_sweep_ts"): + response.headers["X-Poll-Age"] = str(int(store.now() - meta["last_sweep_ts"])) return Overview( healthy=all_healthy and errors == 0, diff --git a/services/cache.py b/services/cache.py index 2b922b6..de56b4e 100644 --- a/services/cache.py +++ b/services/cache.py @@ -38,6 +38,11 @@ def redis_available() -> bool: return _redis is not None +def get_client() -> "redis.Redis | None": + """Return the raw Redis client for primitives not covered by get/set.""" + return _redis + + async def cache_get(key: str) -> Any | None: """GET key from Redis, return deserialized value or None on miss/error.""" if _redis is None: diff --git a/services/enclosure.py b/services/enclosure.py index 418b5f7..ca5beb5 100644 --- a/services/enclosure.py +++ b/services/enclosure.py @@ -4,6 +4,8 @@ import os import re from pathlib import Path +from services.hwgate import gate + logger = logging.getLogger(__name__) ENCLOSURE_BASE = Path("/sys/class/enclosure") @@ -153,14 +155,16 @@ def _parse_slot_number(entry: Path) -> int | None: async def _run_sg_ses(sg_device: str, page: str) -> str | None: - """Run sg_ses for a given page, returning decoded stdout or None.""" + """Run sg_ses for a given page, returning decoded stdout or None. + Hardware-gated.""" try: - proc = await asyncio.create_subprocess_exec( - "sg_ses", f"--page={page}", sg_device, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "sg_ses", f"--page={page}", sg_device, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, stderr = await proc.communicate() if proc.returncode != 0: logger.warning( "sg_ses page %s failed for %s: %s", @@ -177,22 +181,18 @@ async def _run_sg_ses(sg_device: str, page: str) -> str | None: async def get_enclosure_status(sg_device: str) -> dict | None: - """Run sg_ses status (0x02) + element descriptor (0x07) pages and parse. + """Run sg_ses status page (0x02) and parse enclosure health data. - The element descriptor page provides human-readable names for each - element (e.g. "Temp Inlet", "Fan 1"); these are merged into the status - elements so callers get labelled sensors. Descriptor lookup is - best-effort — if page 0x07 is unsupported, elements fall back to bare - indices. + Deliberately only the status page: the element descriptor page (0x07) + is omitted because it errored on the IOM6/Xyratex expanders here + (``couldn't read config page, res=35``) and returned empty names + anyway. Elements fall back to generic labels. If descriptor names are + ever wanted, fetch 0x07 ONCE at startup and cache — never per poll. """ - status_text, desc_text = await asyncio.gather( - _run_sg_ses(sg_device, "0x02"), - _run_sg_ses(sg_device, "0x07"), - ) + status_text = await _run_sg_ses(sg_device, "0x02") if status_text is None: return None - names = _parse_element_descriptors(desc_text) if desc_text else {} - return _parse_ses_page02(status_text, names) + return _parse_ses_page02(status_text) def _norm_type(element_type: str) -> str: diff --git a/services/health.py b/services/health.py new file mode 100644 index 0000000..0a28e3b --- /dev/null +++ b/services/health.py @@ -0,0 +1,21 @@ +"""Shared drive-health classification. + +Single source of truth for turning a parsed SMART dict into a +healthy/warning/error status, used by both the poller (host drives) and +the API (enclosure drives) so the rule never drifts between them. +""" + + +def compute_health_status(smart: dict) -> str: + """Return "healthy", "warning", or "error" for a SMART dict.""" + healthy = smart.get("smart_healthy") + realloc = smart.get("reallocated_sectors") or 0 + pending = smart.get("pending_sectors") or 0 + unc = smart.get("uncorrectable_errors") or 0 + supported = smart.get("smart_supported", True) + + if healthy is False: + return "error" + if realloc > 0 or pending > 0 or unc > 0 or (healthy is None and supported): + return "warning" + return "healthy" diff --git a/services/host.py b/services/host.py index 6938aaa..585fd68 100644 --- a/services/host.py +++ b/services/host.py @@ -3,22 +3,28 @@ import json import logging from services.enclosure import discover_enclosures, list_slots -from services.smart import get_smart_data, scan_megaraid_drives -from services.zfs import get_zfs_pool_map +from services.health import compute_health_status +from services.hwgate import gate +from services.smart import _run_smartctl, scan_megaraid_drives logger = logging.getLogger(__name__) -async def get_host_drives() -> list[dict]: - """Discover non-enclosure block devices and return SMART data for each.""" +async def get_host_drives(pool_map: dict) -> list[dict]: + """Discover non-enclosure block devices and return SMART data for each. + + Poller-side producer: serialized through the hardware gate. ``pool_map`` + is passed in (computed once per sweep) to avoid a second zpool call. + """ # Get all block devices via lsblk try: - proc = await asyncio.create_subprocess_exec( - "lsblk", "-d", "-o", "NAME,SIZE,TYPE,MODEL,ROTA,TRAN", "-J", - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, _ = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "lsblk", "-d", "-o", "NAME,SIZE,TYPE,MODEL,ROTA,TRAN", "-J", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, _ = await proc.communicate() lsblk_data = json.loads(stdout) except (FileNotFoundError, json.JSONDecodeError) as e: logger.warning("lsblk failed: %s", e) @@ -59,10 +65,11 @@ async def get_host_drives() -> list[dict]: host_devices.append({"name": name, "drive_type": drive_type}) - # Fetch SMART + ZFS data concurrently - pool_map = await get_zfs_pool_map() - smart_tasks = [get_smart_data(d["name"]) for d in host_devices] - smart_results = await asyncio.gather(*smart_tasks, return_exceptions=True) + # Fetch SMART for each host disk (serialized by the gate inside _run_smartctl) + smart_results = await asyncio.gather( + *[_run_smartctl(d["name"]) for d in host_devices], + return_exceptions=True, + ) results: list[dict] = [] for dev_info, smart_result in zip(host_devices, smart_results): @@ -72,21 +79,9 @@ async def get_host_drives() -> list[dict]: logger.warning("SMART query failed for host drive %s: %s", name, smart_result) smart = {"device": name, "smart_supported": False} else: - smart, _ = smart_result - - # Compute health_status (same logic as overview.py) - healthy = smart.get("smart_healthy") - realloc = smart.get("reallocated_sectors") or 0 - pending = smart.get("pending_sectors") or 0 - unc = smart.get("uncorrectable_errors") or 0 - - if healthy is False: - health_status = "error" - elif realloc > 0 or pending > 0 or unc > 0 or (healthy is None and smart.get("smart_supported", True)): - health_status = "warning" - else: - health_status = "healthy" + smart = smart_result + health_status = compute_health_status(smart) zfs_info = pool_map.get(name, {}) results.append({ @@ -97,7 +92,7 @@ async def get_host_drives() -> list[dict]: "wwn": smart.get("wwn"), "firmware": smart.get("firmware"), "capacity_bytes": smart.get("capacity_bytes"), - "smart_healthy": healthy, + "smart_healthy": smart.get("smart_healthy"), "smart_supported": smart.get("smart_supported", True), "temperature_c": smart.get("temperature_c"), "power_on_hours": smart.get("power_on_hours"), @@ -116,16 +111,7 @@ async def get_host_drives() -> list[dict]: if has_raid: megaraid_drives = await scan_megaraid_drives() for pd in megaraid_drives: - pd_healthy = pd.get("smart_healthy") - pd_realloc = pd.get("reallocated_sectors") or 0 - pd_pending = pd.get("pending_sectors") or 0 - pd_unc = pd.get("uncorrectable_errors") or 0 - if pd_healthy is False: - pd["health_status"] = "error" - elif pd_realloc > 0 or pd_pending > 0 or pd_unc > 0: - pd["health_status"] = "warning" - else: - pd["health_status"] = "healthy" + pd["health_status"] = compute_health_status(pd) pd["drive_type"] = "physical" pd["physical_drives"] = [] diff --git a/services/hwgate.py b/services/hwgate.py new file mode 100644 index 0000000..8547c75 --- /dev/null +++ b/services/hwgate.py @@ -0,0 +1,27 @@ +"""Global hardware-access gate. + +Every command that touches the SAS bus — smartctl, sg_ses, zpool, ledctl — +must run inside this gate. It serializes hardware I/O so the poller can +never fire a thundering herd of subprocesses at fragile SAS expanders +(which on some shelves drop out and suspend pools under concurrent load). + +Concurrency defaults to 1 (fully serialized). It can be raised via +POLL_CONCURRENCY for healthier backplanes, but 1 is the safe default. + +The gate is a per-process asyncio.Semaphore — only the poller process +calls hardware, so a per-process gate is sufficient. The semaphore is +created lazily on first use so it binds to the running event loop. +""" +import asyncio +import os + +_sem: asyncio.Semaphore | None = None + + +def gate() -> asyncio.Semaphore: + """Return the shared hardware semaphore (use as `async with gate():`).""" + global _sem + if _sem is None: + n = max(1, int(os.environ.get("POLL_CONCURRENCY", "1"))) + _sem = asyncio.Semaphore(n) + return _sem diff --git a/services/leds.py b/services/leds.py index 7bb886b..a9955bd 100644 --- a/services/leds.py +++ b/services/leds.py @@ -1,9 +1,12 @@ import asyncio import re +from services.hwgate import gate -async def set_led(device: str, state: str) -> dict: - """Set the locate LED on a drive via ledctl. + +async def run_led(device: str, state: str) -> dict: + """Set the locate LED on a drive via ledctl. Hardware-gated; runs in the + poller, fed by the Redis LED queue. The API enqueues, never calls this. Args: device: Block device name (e.g. "sda"). Must be alphanumeric. @@ -14,13 +17,14 @@ async def set_led(device: str, state: str) -> dict: if state not in ("locate", "off"): raise ValueError(f"Invalid state: {state}") - proc = await asyncio.create_subprocess_exec( - "ledctl", - f"{state}=/dev/{device}", - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "ledctl", + f"{state}=/dev/{device}", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, stderr = await proc.communicate() if proc.returncode != 0: err = stderr.decode().strip() or stdout.decode().strip() diff --git a/services/smart.py b/services/smart.py index 034d45a..13f0aa3 100644 --- a/services/smart.py +++ b/services/smart.py @@ -5,7 +5,7 @@ import os import re import shutil -from services.cache import cache_get, cache_set +from services.hwgate import gate logger = logging.getLogger(__name__) @@ -29,35 +29,18 @@ def sg_ses_available() -> bool: return shutil.which("sg_ses") is not None -async def get_smart_data(device: str) -> tuple[dict, bool]: - """Run smartctl -a -j against a device, with caching. - - Returns (data, cache_hit) tuple. - """ - # Sanitize device name: only allow alphanumeric and hyphens - if not re.match(r"^[a-zA-Z0-9\-]+$", device): - raise ValueError(f"Invalid device name: {device}") - - cached = await cache_get(f"jbod:smart:{device}") - if cached is not None: - return (cached, True) - - result = await _run_smartctl(device) - await cache_set(f"jbod:smart:{device}", result, SMART_CACHE_TTL) - return (result, False) - - async def _run_smartctl(device: str) -> dict: - """Execute smartctl and parse JSON output.""" + """Execute smartctl and parse JSON output. Hardware-gated.""" dev_path = f"/dev/{device}" try: - proc = await asyncio.create_subprocess_exec( - "smartctl", "-a", "-j", dev_path, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, stderr = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "smartctl", "-a", "-j", dev_path, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, stderr = await proc.communicate() except FileNotFoundError: return {"error": "smartctl not found", "smart_supported": False} @@ -154,12 +137,13 @@ def _parse_smart_json(device: str, data: dict) -> dict: async def scan_megaraid_drives() -> list[dict]: """Discover physical drives behind MegaRAID controllers via smartctl --scan.""" try: - proc = await asyncio.create_subprocess_exec( - "smartctl", "--scan", "-j", - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, _ = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "smartctl", "--scan", "-j", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, _ = await proc.communicate() scan_data = json.loads(stdout) except (FileNotFoundError, json.JSONDecodeError) as e: logger.warning("smartctl --scan failed: %s", e) @@ -179,12 +163,13 @@ async def scan_megaraid_drives() -> list[dict]: dev_path = entry["name"] dev_type = entry["type"] try: - proc = await asyncio.create_subprocess_exec( - "smartctl", "-a", "-j", "-d", dev_type, dev_path, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, _ = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + "smartctl", "-a", "-j", "-d", dev_type, dev_path, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, _ = await proc.communicate() if not stdout: return None data = json.loads(stdout) diff --git a/services/store.py b/services/store.py new file mode 100644 index 0000000..0c5046f --- /dev/null +++ b/services/store.py @@ -0,0 +1,134 @@ +"""Redis as the sole interface between the poller (producer) and the API/ +MQTT consumers. + +The poller writes hardware-derived state here; the API and MQTT publisher +read it and never touch hardware. Also carries the locate-LED command queue +(API enqueues, poller drains) and the poller single-instance lock. +""" +import json +import logging +import os +import time + +from services.cache import cache_get, cache_set, get_client + +logger = logging.getLogger(__name__) + +# Data TTL: must comfortably exceed the sweep interval so values persist +# between sweeps (last-good is served stale rather than blanked). +DATA_TTL = int(os.environ.get("POLL_DATA_TTL", "900")) + +K_SMART = "jbod:smart:{}" +K_SES = "jbod:ses:{}" +K_ZFS = "jbod:zfs_map" +K_INVENTORY = "jbod:inventory" +K_HOST_DRIVES = "jbod:host_drives" +K_META = "jbod:poll:meta" +K_LED_QUEUE = "jbod:led:queue" +K_LOCK = "jbod:poll:lock" + + +# ── writes (poller) ───────────────────────────────────────────────────── +async def set_smart(device: str, data: dict) -> None: + await cache_set(K_SMART.format(device), data, DATA_TTL) + + +async def set_ses(enc_id: str, data: dict) -> None: + await cache_set(K_SES.format(enc_id), data, DATA_TTL) + + +async def set_zfs_map(pool_map: dict) -> None: + await cache_set(K_ZFS, pool_map, DATA_TTL) + + +async def set_inventory(inventory: dict) -> None: + await cache_set(K_INVENTORY, inventory, DATA_TTL) + + +async def set_host_drives(host_drives: list) -> None: + await cache_set(K_HOST_DRIVES, host_drives, DATA_TTL) + + +async def set_meta(meta: dict) -> None: + # Meta has no TTL — it's the heartbeat; staleness is judged from its ts. + await cache_set(K_META, meta, 0) + + +# ── reads (consumers) ─────────────────────────────────────────────────── +async def get_smart(device: str) -> dict | None: + return await cache_get(K_SMART.format(device)) + + +async def get_ses(enc_id: str) -> dict | None: + return await cache_get(K_SES.format(enc_id)) + + +async def get_zfs_map() -> dict: + return await cache_get(K_ZFS) or {} + + +async def get_inventory() -> dict: + return await cache_get(K_INVENTORY) or {"enclosures": []} + + +async def get_host_drives() -> list: + return await cache_get(K_HOST_DRIVES) or [] + + +async def get_meta() -> dict | None: + return await cache_get(K_META) + + +# ── locate-LED command queue ──────────────────────────────────────────── +async def enqueue_led(device: str, state: str) -> bool: + """API side: push a LED request for the poller to execute. Returns False + if Redis is unavailable (so the API can surface 503).""" + client = get_client() + if client is None: + return False + await client.rpush(K_LED_QUEUE, json.dumps({"device": device, "state": state})) + return True + + +async def pop_led(timeout: int = 5) -> dict | None: + """Poller side: block up to `timeout`s for the next LED request.""" + client = get_client() + if client is None: + return None + item = await client.blpop(K_LED_QUEUE, timeout=timeout) + if not item: + return None + _key, raw = item + try: + return json.loads(raw) + except (ValueError, TypeError): + return None + + +# ── poller single-instance lock ───────────────────────────────────────── +async def acquire_lock(token: str, ttl: int = 30) -> bool: + """Try to claim the poller lock. Returns True if acquired/owned.""" + client = get_client() + if client is None: + return True # no Redis → no coordination possible; run anyway + got = await client.set(K_LOCK, token, nx=True, ex=ttl) + if got: + return True + current = await client.get(K_LOCK) + return current == token + + +async def refresh_lock(token: str, ttl: int = 30) -> bool: + """Extend the lock if we still own it.""" + client = get_client() + if client is None: + return True + current = await client.get(K_LOCK) + if current == token: + await client.expire(K_LOCK, ttl) + return True + return False + + +def now() -> float: + return time.time() diff --git a/services/temps.py b/services/temps.py index eb422bd..8a3d9e8 100644 --- a/services/temps.py +++ b/services/temps.py @@ -1,74 +1,57 @@ -"""Per-enclosure temperature aggregation. +"""Per-enclosure temperature aggregation (read-only consumer). -Builds a compact snapshot of enclosure temperatures suitable for Home -Assistant: each named SES temperature sensor plus a derived per-enclosure -hotspot (the hottest reading across SES sensors *and* the drives housed in -that enclosure — drive temps are usually the real hotspot signal). +Reads the poller's data out of Redis — never touches hardware. Each named +SES temperature sensor plus a derived per-enclosure hotspot (the hottest +reading across SES sensors *and* the drives housed in that enclosure — +drive temps are usually the real hotspot signal). """ -import asyncio import logging +import os import socket -from services.enclosure import discover_enclosures, get_enclosure_status, list_slots -from services.smart import get_smart_data +from services import store logger = logging.getLogger(__name__) def get_node_id() -> str: """Stable identifier for this monitor instance (defaults to hostname).""" - import os - return os.environ.get("MQTT_NODE_ID") or socket.gethostname() or "jbod-monitor" async def build_temp_snapshot() -> dict: - """Collect per-enclosure temperature metrics for all enclosures.""" - enclosures = discover_enclosures() - - async def _health(enc): - if enc.get("sg_device"): - return await get_enclosure_status(enc["sg_device"]) - return None - - health_results = await asyncio.gather( - *[_health(enc) for enc in enclosures], return_exceptions=True - ) - + """Collect per-enclosure temperature metrics from the store.""" + inventory = await store.get_inventory() out_enclosures: list[dict] = [] - for enc, health in zip(enclosures, health_results): - if isinstance(health, Exception): - logger.warning("SES health failed for %s: %s", enc["id"], health) - health = None + + for enc in inventory.get("enclosures", []): + ses = await store.get_ses(enc["id"]) or {} sensors: list[dict] = [] sensor_temps: list[float] = [] - if isinstance(health, dict): - for t in health.get("temps", []): - temp = t.get("temperature_c") - sensors.append({ - "index": t["index"], - "name": t.get("name"), - "status": t.get("status", "Unknown"), - "temperature_c": temp, - }) - if isinstance(temp, (int, float)) and temp > 0: - sensor_temps.append(float(temp)) + for t in ses.get("temps", []): + temp = t.get("temperature_c") + sensors.append({ + "index": t["index"], + "name": t.get("name"), + "status": t.get("status", "Unknown"), + "temperature_c": temp, + }) + if isinstance(temp, (int, float)) and temp > 0: + sensor_temps.append(float(temp)) # Drive temperatures for drives housed in this enclosure. drive_temps: list[float] = [] - devices = [s["device"] for s in list_slots(enc["id"]) if s.get("device")] - if devices: - smart_results = await asyncio.gather( - *[get_smart_data(d) for d in devices], return_exceptions=True - ) - for res in smart_results: - if isinstance(res, Exception): - continue - data, _hit = res - temp = data.get("temperature_c") - if isinstance(temp, (int, float)) and temp > 0: - drive_temps.append(float(temp)) + for slot in enc.get("slots", []): + dev = slot.get("device") + if not dev: + continue + sm = await store.get_smart(dev) + if not sm: + continue + temp = sm.get("temperature_c") + if isinstance(temp, (int, float)) and temp > 0: + drive_temps.append(float(temp)) all_temps = sensor_temps + drive_temps out_enclosures.append({ diff --git a/services/zfs.py b/services/zfs.py index 71fcf0a..d2f1b8a 100644 --- a/services/zfs.py +++ b/services/zfs.py @@ -4,26 +4,23 @@ import logging import re from pathlib import Path -from services.cache import cache_get, cache_set +from services.hwgate import gate logger = logging.getLogger(__name__) # Allow overriding the zpool binary path via env (for bind-mounted host tools) ZPOOL_BIN = os.environ.get("ZPOOL_BIN", "zpool") -ZFS_CACHE_TTL = 300 - async def get_zfs_pool_map() -> dict[str, dict]: """Return a dict mapping device names to ZFS pool and vdev info. e.g. {"sda": {"pool": "tank", "vdev": "raidz2-0"}, "sdb": {"pool": "fast", "vdev": "mirror-0"}} - """ - cached = await cache_get("jbod:zfs_map") - if cached is not None: - return cached + Hardware-gated producer; the poller persists the result to the store + and consumers read it from there. + """ pool_map = {} try: # When running in a container with pid:host, use nsenter to run @@ -34,12 +31,13 @@ async def get_zfs_pool_map() -> dict[str, dict]: else: cmd = [ZPOOL_BIN, "status", "-P"] - proc = await asyncio.create_subprocess_exec( - *cmd, - stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, - ) - stdout, _ = await proc.communicate() + async with gate(): + proc = await asyncio.create_subprocess_exec( + *cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, _ = await proc.communicate() if proc.returncode != 0: return pool_map @@ -103,7 +101,6 @@ async def get_zfs_pool_map() -> dict[str, dict]: except FileNotFoundError: logger.debug("zpool not available") - await cache_set("jbod:zfs_map", pool_map, ZFS_CACHE_TTL) return pool_map