From 0f829e03806c788c2a8c7661c85492fd2415ca27 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 15:12:34 +0000 Subject: [PATCH 01/13] 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 -- 2.49.1 From cccc2bc0045cafa551bfacdb4edc165fe98a78b5 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 15:23:42 +0000 Subject: [PATCH 02/13] Split deploy into independent poller and web stacks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Decouple the privileged hardware poller from web iterations so a web redeploy can never recreate/restart the poller (the backplane-touching container). Two images from one Dockerfile via build targets: - jbod-poller: privileged producer, ships smartmontools/sg3-utils/ledmon, Redis-only deps (~197MB) - jbod-web: unprivileged read-only consumer (REST/UI/MQTT), no hardware tools, no /dev (~147MB) Two compose projects sharing Redis over host networking: - compose.infra.yml (jbod-infra): poller + redis — stable substrate - compose.web.yml (jbod-web): app — 'up -d --build' touches only app build.sh builds/pushes both images; split requirements-{poller,web}.txt; dropped the combined docker-compose.yml; .dockerignore excludes tmp/ (keeps the venv and the mqtt.env creds out of the build context). --- .dockerignore | 2 ++ Dockerfile | 29 ++++++++++++----- README.md | 15 +++++++-- build.sh | 23 ++++++++------ compose.infra.yml | 47 +++++++++++++++++++++++++++ compose.web.yml | 33 +++++++++++++++++++ docker-compose.yml | 70 ----------------------------------------- requirements-poller.txt | 2 ++ requirements-web.txt | 6 ++++ 9 files changed, 137 insertions(+), 90 deletions(-) create mode 100644 compose.infra.yml create mode 100644 compose.web.yml delete mode 100644 docker-compose.yml create mode 100644 requirements-poller.txt create mode 100644 requirements-web.txt diff --git a/.dockerignore b/.dockerignore index df7066d..bfd1096 100644 --- a/.dockerignore +++ b/.dockerignore @@ -6,3 +6,5 @@ __pycache__ .git .env *.md +tmp/ +.venv/ diff --git a/Dockerfile b/Dockerfile index 71cf3cc..ef11d8d 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,4 +1,4 @@ -# ---- Stage 1: Build frontend ---- +# ---- Stage: build frontend ---- FROM node:22-alpine AS frontend-build WORKDIR /app/frontend COPY frontend/package.json frontend/package-lock.json* ./ @@ -7,8 +7,10 @@ COPY frontend/ ./ RUN npm run build # Output: /app/static/ -# ---- Stage 2: Production runtime ---- -FROM python:3.13-slim +# ---- Target: poller (hardware producer) ---- +# Privileged at runtime; the only image with SAS/SMART tooling. Slim Python +# deps (just Redis) — no web stack. +FROM python:3.13-slim AS poller RUN apt-get update && apt-get install -y --no-install-recommends \ smartmontools \ @@ -20,16 +22,27 @@ RUN apt-get update && apt-get install -y --no-install-recommends \ WORKDIR /app -COPY requirements.txt . -RUN pip install --no-cache-dir -r requirements.txt +COPY requirements-poller.txt . +RUN pip install --no-cache-dir -r requirements-poller.txt + +COPY poller.py . +COPY services/ services/ + +CMD ["python", "-u", "poller.py"] + +# ---- Target: web (read-only consumer) ---- +# Unprivileged at runtime; no hardware tools. Serves REST/UI + MQTT, reads Redis. +FROM python:3.13-slim AS web + +WORKDIR /app + +COPY requirements-web.txt . +RUN pip install --no-cache-dir -r requirements-web.txt COPY main.py . -COPY poller.py . COPY routers/ routers/ COPY services/ services/ COPY models/ models/ - -# Copy built frontend from stage 1 COPY --from=frontend-build /app/static/ static/ EXPOSE 8000 diff --git a/README.md b/README.md index 0733b4f..b848881 100644 --- a/README.md +++ b/README.md @@ -45,12 +45,21 @@ apt install smartmontools sg3-utils ledmon # Debian/Ubuntu ## Run (docker compose) +Two **independent** stacks so web iterations never recreate the privileged +poller (which touches the SAS bus). They're built as two images — +`jbod-poller` (privileged, has the SAS/SMART tooling) and `jbod-web` (slim, +unprivileged) — and share Redis over host networking. + ```bash -docker compose up -d --build +# 1. Infra: poller + redis — the stable substrate. Deploy once; touch deliberately. +docker compose -f compose.infra.yml up -d + +# 2. Web: REST/UI/MQTT — iterate freely; this only ever recreates `app`. +docker compose -f compose.web.yml up -d --build ``` -Brings up `poller` (privileged), `app` (port 8000), and `redis`. To run the two -processes by hand instead: +Pin both images to a known-good SHA in deploy (`./build.sh` tags `:` and +`:latest` for both). To run the two processes by hand instead: ```bash REDIS_HOST=localhost sudo -E python poller.py # producer (needs root) diff --git a/build.sh b/build.sh index 2abb9cc..0725d9e 100755 --- a/build.sh +++ b/build.sh @@ -1,7 +1,7 @@ #!/usr/bin/env bash set -euo pipefail -IMAGE="docker.adamksmith.xyz/jbod-monitor" +REG="docker.adamksmith.xyz" cd "$(dirname "$0")" @@ -10,16 +10,21 @@ git add -A if git diff --cached --quiet; then echo "No changes to commit, using HEAD" else - git commit -m "${1:-Build and push image}" + git commit -m "${1:-Build and push images}" fi SHA=$(git rev-parse --short HEAD) -echo "Building ${IMAGE}:${SHA}" -docker build -t "${IMAGE}:${SHA}" -t "${IMAGE}:latest" . +# Two images from one Dockerfile (shared frontend build stage is cached): +# jbod-poller — privileged hardware producer +# jbod-web — unprivileged read-only consumer (REST/UI/MQTT) +for target in poller web; do + img="${REG}/jbod-${target}" + echo "Building ${img}:${SHA}" + docker build --target "${target}" -t "${img}:${SHA}" -t "${img}:latest" . + echo "Pushing ${img}:${SHA}" + docker push "${img}:${SHA}" + docker push "${img}:latest" +done -echo "Pushing ${IMAGE}:${SHA}" -docker push "${IMAGE}:${SHA}" -docker push "${IMAGE}:latest" - -echo "Done: ${IMAGE}:${SHA}" +echo "Done: ${REG}/jbod-{poller,web}:${SHA}" diff --git a/compose.infra.yml b/compose.infra.yml new file mode 100644 index 0000000..a2eb852 --- /dev/null +++ b/compose.infra.yml @@ -0,0 +1,47 @@ +# Infra stack: the hardware poller + Redis. This is the STABLE substrate — +# deploy once, pin to a known-good image SHA, and touch it deliberately +# (never as part of a web iteration). Web redeploys use compose.web.yml and +# can never name these services. +# +# docker compose -f compose.infra.yml up -d # bring up / update poller+redis +# +name: jbod-infra + +services: + poller: + image: docker.adamksmith.xyz/jbod-poller:latest # pin to a SHA in deploy + container_name: jbod-poller + restart: unless-stopped + privileged: true + pid: host + network_mode: host + volumes: + - /dev:/dev + - /sys:/sys:ro + - /run/udev:/run/udev:ro + environment: + - TZ=America/Denver + - ZFS_USE_NSENTER=true + - REDIS_HOST=127.0.0.1 + - REDIS_PORT=6379 + - REDIS_DB=0 + # 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 + + redis: + image: redis:7-alpine + container_name: jbod-redis + restart: unless-stopped + network_mode: host + volumes: + - redis-data:/data + command: redis-server --save 60 1 --loglevel warning + +volumes: + redis-data: diff --git a/compose.web.yml b/compose.web.yml new file mode 100644 index 0000000..dd58d64 --- /dev/null +++ b/compose.web.yml @@ -0,0 +1,33 @@ +# Web stack: REST API + UI + Home Assistant MQTT. Read-only consumer of the +# Redis filled by compose.infra.yml. Iterate freely here — this command only +# ever touches the `app` container, never the poller: +# +# docker compose -f compose.web.yml up -d --build +# +# Requires the infra stack (Redis) to be running. Reaches Redis at +# 127.0.0.1:6379 via host networking (separate compose project, shared host net). +name: jbod-web + +services: + app: + image: docker.adamksmith.xyz/jbod-web:latest # pin to a SHA in deploy + 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 (MQTT_USERNAME/MQTT_PASSWORD). Optional: container + # starts without it (MQTT just won't authenticate). + env_file: + - path: /etc/jbod-monitor/mqtt.env + required: false diff --git a/docker-compose.yml b/docker-compose.yml deleted file mode 100644 index 4b1f81e..0000000 --- a/docker-compose.yml +++ /dev/null @@ -1,70 +0,0 @@ -services: - # 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-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 - environment: - - TZ=America/Denver - - ZFS_USE_NSENTER=true - - REDIS_HOST=127.0.0.1 - - REDIS_PORT=6379 - - REDIS_DB=0 - # 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 (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 - - redis: - image: redis:7-alpine - container_name: jbod-redis - restart: unless-stopped - network_mode: host - volumes: - - redis-data:/data - command: redis-server --save 60 1 --loglevel warning - -volumes: - redis-data: diff --git a/requirements-poller.txt b/requirements-poller.txt new file mode 100644 index 0000000..a0390dd --- /dev/null +++ b/requirements-poller.txt @@ -0,0 +1,2 @@ +# Poller image — hardware producer. Only needs Redis; no web stack. +redis>=5.0.0 diff --git a/requirements-web.txt b/requirements-web.txt new file mode 100644 index 0000000..20eeae1 --- /dev/null +++ b/requirements-web.txt @@ -0,0 +1,6 @@ +# Web/API image — read-only consumer + MQTT. No hardware tools needed. +fastapi>=0.115.0 +uvicorn>=0.34.0 +pydantic>=2.10.0 +redis>=5.0.0 +paho-mqtt>=2.1.0 -- 2.49.1 From 7878e17e550d8a7a3aee1970a2a1210293c08c64 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 15:58:43 +0000 Subject: [PATCH 03/13] Harden poller from Codex review: close residual storm/robustness gaps MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adversarial review (gpt-5.5) of the new architecture surfaced real gaps: Hardware-safety / storm prevention: - Restart-storm guard: on startup, defer the first sweep if a sweep ran < POLL_SWEEP_INTERVAL ago (persisted in Redis), so deploy-cycling can no longer trigger back-to-back full hardware sweeps (poller.py). - Centralize pacing in the hardware gate: POLL_DRIVE_GAP is now held after EVERY hardware op (SMART, SES, host, MegaRAID, ledctl), not just the enclosure-drive loop (hwgate.py); removed the now-redundant per-loop sleeps. - Single-instance lock made atomic (Lua compare-and-set / compare-and-expire) and the refresher is now fatal on any error + has a done-callback, so a dropped lock can never leave two pollers sweeping concurrently (store.py, poller.py). - Warn loudly when POLL_CONCURRENCY > 1. Correctness: - Heartbeat now actually writes: cache_set omits EX when ttl<=0 (Redis rejects EX 0), so poller meta/last_sweep_ts persists — this also enables the restart-storm guard and the poller_fresh health flag (cache.py). - Web image runs a single uvicorn worker so the MQTT publisher is a true singleton (two workers shared a client_id and flapped the broker) (Dockerfile). - LED enqueue catches Redis errors -> API returns 503, not 500 (store.py). Deploy independence: - build.sh takes a target (web|poller|all) so web iterations never rebuild or repush the poller image. Nits: drop dead SMART_CACHE_TTL constant. --- Dockerfile | 5 ++++- build.sh | 26 ++++++++++++++++++++------ poller.py | 39 +++++++++++++++++++++++++++++++++------ services/cache.py | 7 +++++-- services/hwgate.py | 46 ++++++++++++++++++++++++++++++++++------------ services/smart.py | 2 -- services/store.py | 42 +++++++++++++++++++++++++++--------------- 7 files changed, 123 insertions(+), 44 deletions(-) diff --git a/Dockerfile b/Dockerfile index ef11d8d..c09ea5e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -50,4 +50,7 @@ EXPOSE 8000 HEALTHCHECK --interval=30s --timeout=5s --retries=3 \ CMD python -c "import urllib.request; urllib.request.urlopen('http://localhost:8000/api/health')" || exit 1 -CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "2"] +# Single worker: the MQTT publisher is a per-process singleton (multiple +# workers would connect with the same client_id and flap the broker). The +# API is a low-traffic internal monitor, so one worker is plenty. +CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "1"] diff --git a/build.sh b/build.sh index 0725d9e..30ff9d2 100755 --- a/build.sh +++ b/build.sh @@ -1,24 +1,38 @@ #!/usr/bin/env bash set -euo pipefail +# Usage: ./build.sh [poller|web|all] [commit message] +# web — build/push ONLY the web image (use this for web iterations so the +# poller image is never rebuilt/repushed) +# poller— build/push ONLY the poller image (touch deliberately) +# all — both (default) +# +# Deploy pins images by SHA (compose.*.yml ship :latest as a convenience only). + REG="docker.adamksmith.xyz" cd "$(dirname "$0")" +TARGET="${1:-all}" +case "$TARGET" in + poller) TARGETS="poller" ;; + web) TARGETS="web" ;; + all) TARGETS="poller web" ;; + *) echo "usage: $0 [poller|web|all] [message]" >&2; exit 1 ;; +esac +MSG="${2:-Build and push images}" + # Stage and commit all changes git add -A if git diff --cached --quiet; then echo "No changes to commit, using HEAD" else - git commit -m "${1:-Build and push images}" + git commit -m "${MSG}" fi SHA=$(git rev-parse --short HEAD) -# Two images from one Dockerfile (shared frontend build stage is cached): -# jbod-poller — privileged hardware producer -# jbod-web — unprivileged read-only consumer (REST/UI/MQTT) -for target in poller web; do +for target in ${TARGETS}; do img="${REG}/jbod-${target}" echo "Building ${img}:${SHA}" docker build --target "${target}" -t "${img}:${SHA}" -t "${img}:latest" . @@ -27,4 +41,4 @@ for target in poller web; do docker push "${img}:latest" done -echo "Done: ${REG}/jbod-{poller,web}:${SHA}" +echo "Done (${TARGET}): ${REG}/jbod-{${TARGETS// /,}}:${SHA}" diff --git a/poller.py b/poller.py index 5bff8a5..e107597 100644 --- a/poller.py +++ b/poller.py @@ -35,9 +35,10 @@ logging.basicConfig( 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 +# Pacing between hardware ops (POLL_DRIVE_GAP) is enforced centrally by the +# hardware gate, so every op — SMART, SES, host, MegaRAID, ledctl — is spaced. TOKEN = f"{socket.gethostname()}-{os.getpid()}" @@ -66,16 +67,15 @@ async def sweep() -> None: except Exception as e: errors.append(f"zfs:{e}") - # 3. Per-drive SMART, one at a time with a gap between each. + # 3. Per-drive SMART (gate serializes + paces each call). 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. + # 4. Per-enclosure SES status (gate serializes + paces each call). for enc in enclosures: if not enc.get("sg_device"): continue @@ -85,7 +85,6 @@ async def sweep() -> 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: @@ -132,11 +131,27 @@ async def led_consumer() -> None: async def lock_refresher() -> None: while True: await asyncio.sleep(LOCK_TTL / 3) - if not await store.refresh_lock(TOKEN, LOCK_TTL): + # Any failure to confirm we still hold the lock is fatal — keeping the + # lock alive is what prevents a second poller from sweeping concurrently. + try: + ok = await store.refresh_lock(TOKEN, LOCK_TTL) + except Exception as e: + logger.error("lock refresh error (%s) — exiting to avoid double-polling", e) + os._exit(1) + if not ok: logger.error("lost poller lock — exiting to avoid double-polling") os._exit(1) +def _critical_task_done(task: asyncio.Task) -> None: + """Fail-safe: if a critical task ends for any reason other than a clean + cancel (shutdown), exit so we never run unguarded.""" + if task.cancelled(): + return + logger.error("critical task ended unexpectedly — exiting") + os._exit(1) + + async def main() -> None: await init_cache() while not redis_available(): @@ -162,7 +177,19 @@ async def main() -> None: logger.info("startup jitter %.1fs", jitter) await asyncio.sleep(jitter) + # Restart-storm guard: if a sweep ran recently (persisted in Redis across + # restarts), wait out the remainder of the interval instead of immediately + # re-sweeping. This is the key protection against deploy-cycling storms. + meta = await store.get_meta() + if meta and meta.get("last_sweep_ts"): + since = store.now() - meta["last_sweep_ts"] + if 0 <= since < SWEEP_INTERVAL: + wait = SWEEP_INTERVAL - since + logger.info("last sweep %.0fs ago — deferring first sweep %.0fs", since, wait) + await asyncio.sleep(wait) + refresher = asyncio.create_task(lock_refresher()) + refresher.add_done_callback(_critical_task_done) leds = asyncio.create_task(led_consumer()) try: while True: diff --git a/services/cache.py b/services/cache.py index de56b4e..a7c9970 100644 --- a/services/cache.py +++ b/services/cache.py @@ -58,10 +58,13 @@ async def cache_get(key: str) -> Any | None: async def cache_set(key: str, value: Any, ttl: int = 120) -> None: - """SET key in Redis with expiry, silently catches errors.""" + """SET key in Redis. ttl<=0 stores with no expiry (Redis rejects EX 0).""" if _redis is None: return try: - await _redis.set(key, json.dumps(value), ex=ttl) + if ttl and ttl > 0: + await _redis.set(key, json.dumps(value), ex=ttl) + else: + await _redis.set(key, json.dumps(value)) except Exception as e: logger.warning("Redis SET %s failed: %s", key, e) diff --git a/services/hwgate.py b/services/hwgate.py index 8547c75..67d7cd9 100644 --- a/services/hwgate.py +++ b/services/hwgate.py @@ -1,27 +1,49 @@ -"""Global hardware-access gate. +"""Global hardware-access gate + pacing. -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). +Every command that touches the SAS bus — smartctl, sg_ses, zpool, lsblk, +ledctl — must run inside this gate. It does two things: -Concurrency defaults to 1 (fully serialized). It can be raised via -POLL_CONCURRENCY for healthier backplanes, but 1 is the safe default. + 1. Serializes hardware I/O (semaphore, default concurrency 1) so the + poller can never fire a thundering herd at fragile SAS expanders. + 2. Paces it: after every hardware op it holds the gate for POLL_DRIVE_GAP + seconds, so back-to-back ops (gathered SMART, MegaRAID scans, SES, + ledctl) are always spaced — not just the per-drive loop. -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. +This is a per-process gate; only the poller process calls hardware, so +per-process serialization is sufficient. Created lazily so it binds to the +running event loop. """ import asyncio +import contextlib +import logging import os +logger = logging.getLogger(__name__) + _sem: asyncio.Semaphore | None = None -def gate() -> asyncio.Semaphore: - """Return the shared hardware semaphore (use as `async with gate():`).""" +def _semaphore() -> asyncio.Semaphore: global _sem if _sem is None: n = max(1, int(os.environ.get("POLL_CONCURRENCY", "1"))) + if n > 1: + logger.warning( + "POLL_CONCURRENCY=%d (>1) — hardware ops will run concurrently; " + "only do this on backplanes known to tolerate it", n, + ) _sem = asyncio.Semaphore(n) return _sem + + +@contextlib.asynccontextmanager +async def gate(): + """Acquire the hardware gate for one op, then hold it for the pacing gap.""" + sem = _semaphore() + async with sem: + try: + yield + finally: + gap = float(os.environ.get("POLL_DRIVE_GAP", "0.5")) + if gap > 0: + await asyncio.sleep(gap) diff --git a/services/smart.py b/services/smart.py index 13f0aa3..0051ef1 100644 --- a/services/smart.py +++ b/services/smart.py @@ -18,8 +18,6 @@ ATTR_PENDING = 197 ATTR_UNCORRECTABLE = 198 ATTR_WEAR_LEVELING = 177 # SSD wear leveling -SMART_CACHE_TTL = int(os.environ.get("SMART_CACHE_TTL", "120")) - def smartctl_available() -> bool: return shutil.which("smartctl") is not None diff --git a/services/store.py b/services/store.py index 0c5046f..ee349d9 100644 --- a/services/store.py +++ b/services/store.py @@ -82,11 +82,15 @@ async def get_meta() -> dict | None: # ── 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).""" + if Redis is unavailable or the push fails (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})) + try: + await client.rpush(K_LED_QUEUE, json.dumps({"device": device, "state": state})) + except Exception as e: + logger.warning("LED enqueue failed: %s", e) + return False return True @@ -105,29 +109,37 @@ async def pop_led(timeout: int = 5) -> dict | None: return None -# ── poller single-instance lock ───────────────────────────────────────── +# ── poller single-instance lock (atomic compare-and-set) ──────────────── +# Acquire if free OR already ours, refreshing the TTL — all in one round trip +# so there's no GET/SET race that could let two pollers run concurrently. +_ACQUIRE_LUA = ( + "local v = redis.call('get', KEYS[1]) " + "if v == false or v == ARGV[1] then " + " redis.call('set', KEYS[1], ARGV[1], 'EX', ARGV[2]) return 1 " + "else return 0 end" +) +# Extend only if we still own it (atomic compare-and-expire). +_REFRESH_LUA = ( + "if redis.call('get', KEYS[1]) == ARGV[1] then " + " return redis.call('expire', KEYS[1], ARGV[2]) " + "else return 0 end" +) + + async def acquire_lock(token: str, ttl: int = 30) -> bool: - """Try to claim the poller lock. Returns True if acquired/owned.""" + """Atomically claim the poller lock (or re-own it). True if held.""" 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 + return bool(await client.eval(_ACQUIRE_LUA, 1, K_LOCK, token, ttl)) async def refresh_lock(token: str, ttl: int = 30) -> bool: - """Extend the lock if we still own it.""" + """Atomically extend the lock iff 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 + return bool(await client.eval(_REFRESH_LUA, 1, K_LOCK, token, ttl)) def now() -> float: -- 2.49.1 From 226f73851c108a91b3b9c88aec3f7fda96777595 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 16:26:36 +0000 Subject: [PATCH 04/13] Fix critical lock-expiry window found in Codex second pass MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - poller.py: start lock_refresher() IMMEDIATELY after acquiring the lock, before the startup jitter and restart-storm deferral. Those sleeps can together exceed LOCK_TTL (restart defer is up to POLL_SWEEP_INTERVAL), so with the refresher started afterward the lock could expire mid-deferral and a second poller could begin sweeping concurrently — the exact hardware-storm risk this design prevents. - hwgate.py: clamp POLL_CONCURRENCY to 1 unless POLL_ALLOW_UNSAFE_CONCURRENCY is set (fragile SAS hardware; >1 can drop expanders). Warn either way. - build.sh: keep backward-compatible — a non-target first arg is treated as the commit message (./build.sh "msg" still works), known targets select web|poller|all. --- build.sh | 15 ++++++++------- poller.py | 9 +++++++-- services/hwgate.py | 18 +++++++++++++++--- 3 files changed, 30 insertions(+), 12 deletions(-) diff --git a/build.sh b/build.sh index 30ff9d2..7519654 100755 --- a/build.sh +++ b/build.sh @@ -13,14 +13,15 @@ REG="docker.adamksmith.xyz" cd "$(dirname "$0")" -TARGET="${1:-all}" -case "$TARGET" in - poller) TARGETS="poller" ;; - web) TARGETS="web" ;; - all) TARGETS="poller web" ;; - *) echo "usage: $0 [poller|web|all] [message]" >&2; exit 1 ;; +# First arg is the target if it's a known one; otherwise it's treated as the +# commit message (backward-compatible with `./build.sh "message"`). +case "${1:-}" in + poller) TARGETS="poller"; MSG="${2:-Build and push images}" ;; + web) TARGETS="web"; MSG="${2:-Build and push images}" ;; + all|"") TARGETS="poller web"; MSG="${2:-Build and push images}" ;; + *) TARGETS="poller web"; MSG="${1}" ;; esac -MSG="${2:-Build and push images}" +TARGET="${TARGETS// /+}" # Stage and commit all changes git add -A diff --git a/poller.py b/poller.py index e107597..4ea1820 100644 --- a/poller.py +++ b/poller.py @@ -165,6 +165,13 @@ async def main() -> None: await asyncio.sleep(LOCK_TTL // 2) logger.info("poller lock acquired (%s)", TOKEN) + # Start the refresher IMMEDIATELY — before any pre-sweep sleep. The jitter + # and restart-storm deferral below can together exceed LOCK_TTL, so without + # an active refresher the lock would expire and a second poller could begin + # sweeping concurrently (a hardware-storm risk). + refresher = asyncio.create_task(lock_refresher()) + refresher.add_done_callback(_critical_task_done) + if not smartctl_available(): logger.warning("smartctl not found — install smartmontools") if not sg_ses_available(): @@ -188,8 +195,6 @@ async def main() -> None: logger.info("last sweep %.0fs ago — deferring first sweep %.0fs", since, wait) await asyncio.sleep(wait) - refresher = asyncio.create_task(lock_refresher()) - refresher.add_done_callback(_critical_task_done) leds = asyncio.create_task(led_consumer()) try: while True: diff --git a/services/hwgate.py b/services/hwgate.py index 67d7cd9..d4d54a4 100644 --- a/services/hwgate.py +++ b/services/hwgate.py @@ -27,11 +27,23 @@ def _semaphore() -> asyncio.Semaphore: global _sem if _sem is None: n = max(1, int(os.environ.get("POLL_CONCURRENCY", "1"))) + # Serialized (1) is the safe default. Concurrent hardware I/O can drop + # fragile SAS expanders, so >1 is clamped unless explicitly overridden. if n > 1: - logger.warning( - "POLL_CONCURRENCY=%d (>1) — hardware ops will run concurrently; " - "only do this on backplanes known to tolerate it", n, + override = os.environ.get("POLL_ALLOW_UNSAFE_CONCURRENCY", "").lower() in ( + "1", "true", "yes", ) + if override: + logger.warning( + "POLL_CONCURRENCY=%d with unsafe override — concurrent " + "hardware I/O enabled; backplane must tolerate it", n, + ) + else: + logger.warning( + "POLL_CONCURRENCY=%d ignored (clamped to 1 for hardware " + "safety); set POLL_ALLOW_UNSAFE_CONCURRENCY=1 to override", n, + ) + n = 1 _sem = asyncio.Semaphore(n) return _sem -- 2.49.1 From 44a78244256c1d3d506109afa8949a9d3c75ba6d Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 17:21:16 +0000 Subject: [PATCH 05/13] Name the MQTT hub device (fix HA 'Unnamed Device') Enclosure entities set via_device to a parent hub id that nothing defined, so Home Assistant showed an 'Unnamed Device'. Publish a hub connectivity binary_sensor ('Status', online/offline from the LWT topic) that defines the parent device with a proper name, so the per-enclosure devices nest cleanly under 'JBOD Monitor ()'. --- services/mqtt_publisher.py | 38 ++++++++++++++++++++++++++++++++++++-- 1 file changed, 36 insertions(+), 2 deletions(-) diff --git a/services/mqtt_publisher.py b/services/mqtt_publisher.py index 23c1bd0..c3c9b59 100644 --- a/services/mqtt_publisher.py +++ b/services/mqtt_publisher.py @@ -80,7 +80,10 @@ class MqttPublisher: self.discovery_prefix = os.environ.get("MQTT_DISCOVERY_PREFIX", "homeassistant") self.base_topic = os.environ.get("MQTT_BASE_TOPIC", "jbod-monitor") self.interval = int(os.environ.get("MQTT_PUBLISH_INTERVAL", "60")) - self.node = _slug(get_node_id()) + self.node_raw = get_node_id() + self.node = _slug(self.node_raw) + # Parent "hub" device that the per-enclosure devices nest under. + self.hub_id = f"{self.node}_jbod_monitor" self.availability_topic = f"{self.base_topic}/{self.node}/status" # unique_ids for which we've already published discovery this connection. @@ -127,6 +130,16 @@ class MqttPublisher: logger.warning("MQTT disconnected: %s", reason_code) # discovery / state ------------------------------------------------------ + def _hub_device_block(self) -> dict: + """The parent device the enclosures nest under (named, so HA doesn't + show an 'Unnamed Device').""" + return { + "identifiers": [self.hub_id], + "name": f"JBOD Monitor ({self.node_raw})", + "manufacturer": "jbod-monitor", + "model": "SES/SMART poller", + } + def _device_block(self, enc: dict) -> dict: enc_id = enc["id"] vendor = enc.get("vendor", "").strip() @@ -137,9 +150,29 @@ class MqttPublisher: "name": f"JBOD {name} ({enc_id})", "manufacturer": vendor or None, "model": model or None, - "via_device": f"{self.node}_jbod_monitor", + "via_device": self.hub_id, } + def _publish_hub(self) -> None: + """Announce the hub device via a connectivity binary_sensor, so the + parent device is named and shows monitor online/offline status.""" + uid = f"{self.hub_id}_status" + if uid in self._discovered: + return + payload = { + "name": "Status", + "unique_id": uid, + "device_class": "connectivity", + "state_topic": self.availability_topic, + "payload_on": "online", + "payload_off": "offline", + "entity_category": "diagnostic", + "device": self._hub_device_block(), + } + topic = f"{self.discovery_prefix}/binary_sensor/{self.hub_id}/status/config" + self._client.publish(topic, json.dumps(payload), qos=1, retain=True) + self._discovered.add(uid) + def _publish_discovery(self, enc: dict) -> None: enc_id = enc["id"] enc_slug = _slug(enc_id) @@ -204,6 +237,7 @@ class MqttPublisher: self._client.publish(state_topic, json.dumps(payload), qos=0, retain=True) def publish_snapshot(self, snapshot: dict) -> None: + self._publish_hub() # name the parent device before nesting enclosures for enc in snapshot.get("enclosures", []): self._publish_discovery(enc) self._publish_state(enc) -- 2.49.1 From 1839e9d5f29e4c6ff794035cc657c35adc085717 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 18:27:56 +0000 Subject: [PATCH 06/13] Add host-poller: chassis IPMI/iDRAC + CPU + on-metal drive temps New third producer container (host-poller) owning the server chassis subsystem, separate from the SAS-shelf poller: - services/host_sensors.py: in-band IPMI (ipmitool sdr) for Inlet/Exhaust/etc. + CPU package temps from coretemp hwmon - host_poller.py: own single-instance lock + heartbeat (host_poll keys), restart-storm guard, paced; writes jbod:host_sensors + jbod:host_drives - move host/on-metal drive collection off the SAS poller to host-poller (reads jbod:zfs_map for annotation) - store: generalized lock/meta with a key param; host_sensors + host-poll keys - mqtt_publisher: publish env/CPU/host-drive temps as entities on the hub device (JBOD Monitor) - Dockerfile host-poller target (ipmitool+smartmontools); compose.infra adds host-poller; build.sh gains host-poller/all --- Dockerfile | 21 ++++++ build.sh | 9 ++- compose.infra.yml | 26 +++++++ host_poller.py | 151 +++++++++++++++++++++++++++++++++++++ poller.py | 8 +- services/host_sensors.py | 95 +++++++++++++++++++++++ services/mqtt_publisher.py | 63 +++++++++++++++- services/store.py | 29 ++++--- 8 files changed, 381 insertions(+), 21 deletions(-) create mode 100644 host_poller.py create mode 100644 services/host_sensors.py diff --git a/Dockerfile b/Dockerfile index c09ea5e..ef2c990 100644 --- a/Dockerfile +++ b/Dockerfile @@ -30,6 +30,27 @@ COPY services/ services/ CMD ["python", "-u", "poller.py"] +# ---- Target: host-poller (chassis IPMI/iDRAC + CPU + on-metal drives) ---- +# Privileged at runtime (needs /dev/ipmi0 + raw disks). Carries ipmitool + +# smartmontools; no SAS-shelf tooling. +FROM python:3.13-slim AS host-poller + +RUN apt-get update && apt-get install -y --no-install-recommends \ + ipmitool \ + smartmontools \ + util-linux \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /app + +COPY requirements-poller.txt . +RUN pip install --no-cache-dir -r requirements-poller.txt + +COPY host_poller.py . +COPY services/ services/ + +CMD ["python", "-u", "host_poller.py"] + # ---- Target: web (read-only consumer) ---- # Unprivileged at runtime; no hardware tools. Serves REST/UI + MQTT, reads Redis. FROM python:3.13-slim AS web diff --git a/build.sh b/build.sh index 7519654..0e00e2a 100755 --- a/build.sh +++ b/build.sh @@ -16,10 +16,11 @@ cd "$(dirname "$0")" # First arg is the target if it's a known one; otherwise it's treated as the # commit message (backward-compatible with `./build.sh "message"`). case "${1:-}" in - poller) TARGETS="poller"; MSG="${2:-Build and push images}" ;; - web) TARGETS="web"; MSG="${2:-Build and push images}" ;; - all|"") TARGETS="poller web"; MSG="${2:-Build and push images}" ;; - *) TARGETS="poller web"; MSG="${1}" ;; + poller) TARGETS="poller"; MSG="${2:-Build and push images}" ;; + host-poller) TARGETS="host-poller"; MSG="${2:-Build and push images}" ;; + web) TARGETS="web"; MSG="${2:-Build and push images}" ;; + all|"") TARGETS="poller host-poller web"; MSG="${2:-Build and push images}" ;; + *) TARGETS="poller host-poller web"; MSG="${1}" ;; esac TARGET="${TARGETS// /+}" diff --git a/compose.infra.yml b/compose.infra.yml index a2eb852..0533a23 100644 --- a/compose.infra.yml +++ b/compose.infra.yml @@ -34,6 +34,32 @@ services: depends_on: - redis + # Chassis producer: IPMI/iDRAC + CPU coretemp + on-metal drives → Redis. + # Separate from the SAS poller (different subsystem); writes host_sensors + + # host_drives. The web consumer publishes these on the hub device. + host-poller: + image: docker.adamksmith.xyz/jbod-host-poller:latest # pin to a SHA in deploy + container_name: jbod-host-poller + restart: unless-stopped + privileged: true + network_mode: host + volumes: + - /dev:/dev + - /sys:/sys:ro + - /run/udev:/run/udev:ro + environment: + - TZ=America/Denver + - REDIS_HOST=127.0.0.1 + - REDIS_PORT=6379 + - REDIS_DB=0 + - HOST_POLL_INTERVAL=120 + - POLL_DRIVE_GAP=0.5 + - POLL_CONCURRENCY=1 + - POLL_STARTUP_JITTER=10 + - POLL_DATA_TTL=900 + depends_on: + - redis + redis: image: redis:7-alpine container_name: jbod-redis diff --git a/host_poller.py b/host_poller.py new file mode 100644 index 0000000..ff6a407 --- /dev/null +++ b/host_poller.py @@ -0,0 +1,151 @@ +"""Host-poller — server chassis temps (IPMI/iDRAC + CPU) and on-metal drives. + +Separate process/container from the SAS poller: this owns the chassis +subsystem (in-band IPMI, CPU coretemp, PERC/onboard drives), which is +independent of the external SAS shelves and shouldn't share their hardware +gate's bus. Writes to Redis; the web/MQTT consumer publishes these on the hub +device. Own single-instance lock + heartbeat; paced via the per-process gate. + +Host-drive ZFS annotation reuses jbod:zfs_map written by the SAS poller. +""" +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.host import get_host_drives +from services.host_sensors import read_coretemp, read_ipmi_temps + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", +) +logger = logging.getLogger("host-poller") + +SWEEP_INTERVAL = int(os.environ.get("HOST_POLL_INTERVAL", "120")) +STARTUP_JITTER = float(os.environ.get("POLL_STARTUP_JITTER", "10")) +LOCK_TTL = 30 + +TOKEN = f"{socket.gethostname()}-hostpoll-{os.getpid()}" + + +async def sweep() -> None: + t0 = time.monotonic() + errors: list[str] = [] + + # 1. IPMI (iDRAC) environmental + CPU temps. + try: + ipmi = await read_ipmi_temps() + except Exception as e: + errors.append(f"ipmi:{e}") + ipmi = [] + env = [s for s in ipmi if s.get("kind") == "env"] + # Prefer kernel coretemp for CPU package temps; fall back to IPMI CPU temps. + cpu = read_coretemp() or [s for s in ipmi if s.get("kind") == "cpu"] + + await store.set_host_sensors({ + "node": socket.gethostname(), + "env": env, + "cpu": cpu, + "ts": store.now(), + }) + + # 2. On-metal (non-enclosure) drives — annotated from the SAS poller's zfs map. + try: + pool_map = await store.get_zfs_map() + host_drives = await get_host_drives(pool_map) + await store.set_host_drives(host_drives) + except Exception as e: + errors.append(f"hostdrives:{e}") + host_drives = [] + + duration = round(time.monotonic() - t0, 1) + await store.set_meta({ + "last_sweep_ts": store.now(), + "last_sweep_duration": duration, + "env_count": len(env), + "cpu_count": len(cpu), + "host_drive_count": len(host_drives), + "errors": errors[:20], + "sweep_interval": SWEEP_INTERVAL, + }, key=store.K_HOST_META) + logger.info( + "host sweep done: %d env, %d cpu, %d drives, %.1fs, %d error(s)", + len(env), len(cpu), len(host_drives), duration, len(errors), + ) + + +async def lock_refresher() -> None: + while True: + await asyncio.sleep(LOCK_TTL / 3) + try: + ok = await store.refresh_lock(TOKEN, LOCK_TTL, key=store.K_HOST_LOCK) + except Exception as e: + logger.error("lock refresh error (%s) — exiting", e) + os._exit(1) + if not ok: + logger.error("lost host-poller lock — exiting") + os._exit(1) + + +def _critical_task_done(task: asyncio.Task) -> None: + if task.cancelled(): + return + logger.error("critical task ended unexpectedly — exiting") + os._exit(1) + + +async def main() -> None: + await init_cache() + while not redis_available(): + logger.warning("Redis unavailable — host-poller needs Redis; retrying in 5s") + await asyncio.sleep(5) + await init_cache() + + while not await store.acquire_lock(TOKEN, LOCK_TTL, key=store.K_HOST_LOCK): + logger.warning("another host-poller holds the lock; waiting %ds", LOCK_TTL // 2) + await asyncio.sleep(LOCK_TTL // 2) + logger.info("host-poller lock acquired (%s)", TOKEN) + + refresher = asyncio.create_task(lock_refresher()) + refresher.add_done_callback(_critical_task_done) + + if os.geteuid() != 0: + logger.warning("not running as root — ipmitool/smartctl may fail") + + jitter = random.uniform(0, STARTUP_JITTER) + logger.info("startup jitter %.1fs", jitter) + await asyncio.sleep(jitter) + + # Restart-storm guard (mirrors the SAS poller). + meta = await store.get_meta(key=store.K_HOST_META) + if meta and meta.get("last_sweep_ts"): + since = store.now() - meta["last_sweep_ts"] + if 0 <= since < SWEEP_INTERVAL: + wait = SWEEP_INTERVAL - since + logger.info("last host sweep %.0fs ago — deferring first sweep %.0fs", since, wait) + await asyncio.sleep(wait) + + try: + while True: + t0 = time.monotonic() + try: + await sweep() + except Exception as e: + logger.error("host sweep error: %s", e) + elapsed = time.monotonic() - t0 + await asyncio.sleep(max(5, SWEEP_INTERVAL - elapsed)) + finally: + refresher.cancel() + await close_cache() + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + pass diff --git a/poller.py b/poller.py index 4ea1820..ce242fc 100644 --- a/poller.py +++ b/poller.py @@ -23,7 +23,6 @@ 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 @@ -86,12 +85,7 @@ async def sweep() -> None: except Exception as e: errors.append(f"ses:{enc['id']}:{e}") - # 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}") + # NOTE: host/on-metal drives are owned by the separate host-poller now. # 6. Inventory + heartbeat. await store.set_inventory({"enclosures": inv_enclosures, "ts": store.now()}) diff --git a/services/host_sensors.py b/services/host_sensors.py new file mode 100644 index 0000000..9bdc191 --- /dev/null +++ b/services/host_sensors.py @@ -0,0 +1,95 @@ +"""Host/chassis temperature collectors: IPMI (iDRAC) + CPU coretemp. + +Used by the host-poller (separate process from the SAS poller). In-band IPMI +via ipmitool (/dev/ipmi0) gives environmental sensors (Inlet/Exhaust/…); the +kernel coretemp hwmon gives CPU package temps. IPMI is gated for serialization; +coretemp is a cheap sysfs read. +""" +import asyncio +import logging +import re +from pathlib import Path + +from services.hwgate import gate + +logger = logging.getLogger(__name__) + +HWMON = Path("/sys/class/hwmon") + + +async def read_ipmi_temps() -> list[dict]: + """`ipmitool sdr type temperature` → [{name, temperature_c, status, kind}].""" + try: + async with gate(): + proc = await asyncio.create_subprocess_exec( + "ipmitool", "sdr", "type", "temperature", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + out, err = await proc.communicate() + except FileNotFoundError: + logger.warning("ipmitool not found") + return [] + if proc.returncode != 0: + logger.warning("ipmitool failed: %s", err.decode(errors="replace").strip()) + return [] + return _parse_ipmi(out.decode(errors="replace")) + + +def _parse_ipmi(text: str) -> list[dict]: + # e.g. "Inlet Temp | 04h | ok | 7.1 | 23 degrees C" + # "Temp | 0Eh | ok | 3.1 | 40 degrees C" (CPU1) + sensors: list[dict] = [] + for line in text.splitlines(): + parts = [p.strip() for p in line.split("|")] + if len(parts) < 5: + continue + name, _sid, status, entity, reading = parts[:5] + m = re.search(r"(-?\d+(?:\.\d+)?)\s*degrees?\s*C", reading, re.IGNORECASE) + if not m: # "no reading", "disabled", "ns", etc. + continue + temp = float(m.group(1)) + kind = "env" + # Dell reports CPU temps as "Temp" under entity 3.x — relabel them. + if name.lower() == "temp" and entity.startswith("3."): + name = f"CPU {entity.split('.')[-1]}" + kind = "cpu" + sensors.append({ + "name": name, + "temperature_c": temp, + "status": status or "ok", + "kind": kind, + }) + return sensors + + +def read_coretemp() -> list[dict]: + """CPU package temps from coretemp hwmon → [{name, temperature_c, kind}].""" + cpus: list[dict] = [] + if not HWMON.is_dir(): + return cpus + for hw in sorted(HWMON.iterdir()): + try: + if (hw / "name").read_text().strip() != "coretemp": + continue + except OSError: + continue + for label_f in sorted(hw.glob("temp*_label")): + try: + label = label_f.read_text().strip() + except OSError: + continue + if not label.lower().startswith("package id"): + continue + input_f = label_f.with_name(label_f.name.replace("_label", "_input")) + try: + milli = int(input_f.read_text().strip()) + except (OSError, ValueError): + continue + pkg = label.split("id")[-1].strip() + cpus.append({ + "name": f"CPU {pkg}", + "temperature_c": round(milli / 1000), + "kind": "cpu", + }) + return cpus diff --git a/services/mqtt_publisher.py b/services/mqtt_publisher.py index c3c9b59..b5bf246 100644 --- a/services/mqtt_publisher.py +++ b/services/mqtt_publisher.py @@ -19,6 +19,7 @@ import logging import os import re +from services import store from services.secrets import read_kv2 from services.temps import build_temp_snapshot, get_node_id @@ -242,6 +243,61 @@ class MqttPublisher: self._publish_discovery(enc) self._publish_state(enc) + def publish_host(self, sensors: dict, drives: list) -> None: + """Publish chassis temps (IPMI env, CPU, on-metal drives) on the hub.""" + self._publish_hub() + state_topic = f"{self.base_topic}/{self.node}/host/state" + device = self._hub_device_block() + payload: dict = {} + + def announce(object_id: str, cfg: dict) -> None: + uid = cfg["unique_id"] + if uid not in self._discovered: + topic = f"{self.discovery_prefix}/sensor/{self.hub_id}/{object_id}/config" + self._client.publish(topic, json.dumps(cfg), qos=1, retain=True) + self._discovered.add(uid) + + common = { + "device_class": "temperature", + "unit_of_measurement": "°C", + "state_class": "measurement", + "state_topic": state_topic, + "availability_topic": self.availability_topic, + "device": device, + } + + def add(kind: str, name: str, temp, icon: str | None = None) -> None: + if temp is None or not name: + return + key = f"{kind}_{_slug(name)}" + payload[key] = temp + cfg = { + **common, + "name": name, + "unique_id": f"{self.hub_id}_{key}", + "value_template": f"{{{{ value_json.{key} }}}}", + } + if icon: + cfg["icon"] = icon + announce(key, cfg) + + for s in sensors.get("env", []): + add("env", s.get("name"), s.get("temperature_c")) + for c in sensors.get("cpu", []): + add("cpu", c.get("name"), c.get("temperature_c"), icon="mdi:cpu-64-bit") + + def add_drive(d: dict) -> None: + dev = d.get("device") + if dev: + add("drive", f"Drive {dev}", d.get("temperature_c"), icon="mdi:harddisk") + + for d in drives: + add_drive(d) + for pd in d.get("physical_drives", []): + add_drive(pd) + + self._client.publish(state_topic, json.dumps(payload), qos=0, retain=True) + async def run(self) -> None: """Background loop: build a snapshot and publish on an interval.""" await asyncio.sleep(3) # let the first SMART poll populate the cache @@ -249,9 +305,14 @@ class MqttPublisher: try: snapshot = await build_temp_snapshot() self.publish_snapshot(snapshot) + host_sensors = await store.get_host_sensors() or {} + host_drives = await store.get_host_drives() or [] + self.publish_host(host_sensors, host_drives) logger.info( - "MQTT published temps for %d enclosure(s)", + "MQTT published temps for %d enclosure(s) + host (%d env, %d cpu)", len(snapshot.get("enclosures", [])), + len(host_sensors.get("env", [])), + len(host_sensors.get("cpu", [])), ) except asyncio.CancelledError: raise diff --git a/services/store.py b/services/store.py index ee349d9..6ed0e1a 100644 --- a/services/store.py +++ b/services/store.py @@ -23,9 +23,12 @@ K_SES = "jbod:ses:{}" K_ZFS = "jbod:zfs_map" K_INVENTORY = "jbod:inventory" K_HOST_DRIVES = "jbod:host_drives" +K_HOST_SENSORS = "jbod:host_sensors" K_META = "jbod:poll:meta" +K_HOST_META = "jbod:host_poll:meta" K_LED_QUEUE = "jbod:led:queue" K_LOCK = "jbod:poll:lock" +K_HOST_LOCK = "jbod:host_poll:lock" # ── writes (poller) ───────────────────────────────────────────────────── @@ -49,9 +52,13 @@ 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: +async def set_host_sensors(sensors: dict) -> None: + await cache_set(K_HOST_SENSORS, sensors, DATA_TTL) + + +async def set_meta(meta: dict, key: str = K_META) -> None: # Meta has no TTL — it's the heartbeat; staleness is judged from its ts. - await cache_set(K_META, meta, 0) + await cache_set(key, meta, 0) # ── reads (consumers) ─────────────────────────────────────────────────── @@ -75,8 +82,12 @@ 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) +async def get_host_sensors() -> dict | None: + return await cache_get(K_HOST_SENSORS) + + +async def get_meta(key: str = K_META) -> dict | None: + return await cache_get(key) # ── locate-LED command queue ──────────────────────────────────────────── @@ -126,20 +137,20 @@ _REFRESH_LUA = ( ) -async def acquire_lock(token: str, ttl: int = 30) -> bool: - """Atomically claim the poller lock (or re-own it). True if held.""" +async def acquire_lock(token: str, ttl: int = 30, key: str = K_LOCK) -> bool: + """Atomically claim a single-instance lock (or re-own it). True if held.""" client = get_client() if client is None: return True # no Redis → no coordination possible; run anyway - return bool(await client.eval(_ACQUIRE_LUA, 1, K_LOCK, token, ttl)) + return bool(await client.eval(_ACQUIRE_LUA, 1, key, token, ttl)) -async def refresh_lock(token: str, ttl: int = 30) -> bool: +async def refresh_lock(token: str, ttl: int = 30, key: str = K_LOCK) -> bool: """Atomically extend the lock iff we still own it.""" client = get_client() if client is None: return True - return bool(await client.eval(_REFRESH_LUA, 1, K_LOCK, token, ttl)) + return bool(await client.eval(_REFRESH_LUA, 1, key, token, ttl)) def now() -> float: -- 2.49.1 From 822829aae4df0d80a862b415773eb01537ee4a21 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 18:45:33 +0000 Subject: [PATCH 07/13] host-poller: poll chassis sensors every 60s (inlet temp 1/min) --- compose.infra.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/compose.infra.yml b/compose.infra.yml index 0533a23..70379c5 100644 --- a/compose.infra.yml +++ b/compose.infra.yml @@ -52,7 +52,7 @@ services: - REDIS_HOST=127.0.0.1 - REDIS_PORT=6379 - REDIS_DB=0 - - HOST_POLL_INTERVAL=120 + - HOST_POLL_INTERVAL=60 - POLL_DRIVE_GAP=0.5 - POLL_CONCURRENCY=1 - POLL_STARTUP_JITTER=10 -- 2.49.1 From ea045433e5373c2274eb74d9a80daaa1d12a207b Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 18:50:16 +0000 Subject: [PATCH 08/13] Split host-poller into separate env (IPMI/CPU) and drive (SMART) loops Two concurrent loops in one process sharing the lock + hardware gate (IPMI and SMART never run concurrently) but on separate schedules: env every HOST_ENV_INTERVAL (60s, responsive inlet/CPU) and drives every HOST_DRIVE_INTERVAL (300s, gentle SMART). Each loop has its own heartbeat (env_meta/drive_meta) + restart-storm guard. --- compose.infra.yml | 3 +- host_poller.py | 121 ++++++++++++++++++++++++++-------------------- services/store.py | 3 +- 3 files changed, 73 insertions(+), 54 deletions(-) diff --git a/compose.infra.yml b/compose.infra.yml index 70379c5..f9edefc 100644 --- a/compose.infra.yml +++ b/compose.infra.yml @@ -52,7 +52,8 @@ services: - REDIS_HOST=127.0.0.1 - REDIS_PORT=6379 - REDIS_DB=0 - - HOST_POLL_INTERVAL=60 + - HOST_ENV_INTERVAL=60 # IPMI/iDRAC + CPU — responsive (inlet temp 1/min) + - HOST_DRIVE_INTERVAL=300 # on-metal SMART — gentle, isolated from env - POLL_DRIVE_GAP=0.5 - POLL_CONCURRENCY=1 - POLL_STARTUP_JITTER=10 diff --git a/host_poller.py b/host_poller.py index ff6a407..f8e1657 100644 --- a/host_poller.py +++ b/host_poller.py @@ -1,12 +1,16 @@ -"""Host-poller — server chassis temps (IPMI/iDRAC + CPU) and on-metal drives. +"""Host-poller — server chassis temps, split into two independent loops. -Separate process/container from the SAS poller: this owns the chassis -subsystem (in-band IPMI, CPU coretemp, PERC/onboard drives), which is -independent of the external SAS shelves and shouldn't share their hardware -gate's bus. Writes to Redis; the web/MQTT consumer publishes these on the hub -device. Own single-instance lock + heartbeat; paced via the per-process gate. +Separate process/container from the SAS poller. Two concurrent loops in one +process, sharing the single-instance lock AND the per-process hardware gate +(so IPMI and SMART never hit hardware at the same instant), but on separate +schedules: -Host-drive ZFS annotation reuses jbod:zfs_map written by the SAS poller. + - env loop (fast): IPMI/iDRAC env + CPU coretemp -> jbod:host_sensors + - drive loop (slow): on-metal drive SMART -> jbod:host_drives + +Keeping SMART on a gentler cadence than the responsive inlet/CPU reads. Each +loop has its own heartbeat + restart-storm guard. Host-drive ZFS annotation +reuses jbod:zfs_map written by the SAS poller. """ import asyncio import logging @@ -26,25 +30,24 @@ logging.basicConfig( ) logger = logging.getLogger("host-poller") -SWEEP_INTERVAL = int(os.environ.get("HOST_POLL_INTERVAL", "120")) +ENV_INTERVAL = int(os.environ.get("HOST_ENV_INTERVAL", "60")) +DRIVE_INTERVAL = int(os.environ.get("HOST_DRIVE_INTERVAL", "300")) STARTUP_JITTER = float(os.environ.get("POLL_STARTUP_JITTER", "10")) LOCK_TTL = 30 TOKEN = f"{socket.gethostname()}-hostpoll-{os.getpid()}" -async def sweep() -> None: +async def sweep_env() -> None: + """IPMI/iDRAC environmental + CPU temps (the responsive loop).""" t0 = time.monotonic() errors: list[str] = [] - - # 1. IPMI (iDRAC) environmental + CPU temps. try: ipmi = await read_ipmi_temps() except Exception as e: errors.append(f"ipmi:{e}") ipmi = [] env = [s for s in ipmi if s.get("kind") == "env"] - # Prefer kernel coretemp for CPU package temps; fall back to IPMI CPU temps. cpu = read_coretemp() or [s for s in ipmi if s.get("kind") == "cpu"] await store.set_host_sensors({ @@ -53,30 +56,57 @@ async def sweep() -> None: "cpu": cpu, "ts": store.now(), }) - - # 2. On-metal (non-enclosure) drives — annotated from the SAS poller's zfs map. - try: - pool_map = await store.get_zfs_map() - host_drives = await get_host_drives(pool_map) - await store.set_host_drives(host_drives) - except Exception as e: - errors.append(f"hostdrives:{e}") - host_drives = [] - - duration = round(time.monotonic() - t0, 1) await store.set_meta({ "last_sweep_ts": store.now(), - "last_sweep_duration": duration, + "last_sweep_duration": round(time.monotonic() - t0, 1), "env_count": len(env), "cpu_count": len(cpu), - "host_drive_count": len(host_drives), "errors": errors[:20], - "sweep_interval": SWEEP_INTERVAL, - }, key=store.K_HOST_META) - logger.info( - "host sweep done: %d env, %d cpu, %d drives, %.1fs, %d error(s)", - len(env), len(cpu), len(host_drives), duration, len(errors), - ) + "interval": ENV_INTERVAL, + }, key=store.K_HOST_ENV_META) + logger.info("env sweep: %d env, %d cpu, %d error(s)", len(env), len(cpu), len(errors)) + + +async def sweep_drives() -> None: + """On-metal (non-enclosure) drive SMART (the gentle loop).""" + t0 = time.monotonic() + errors: list[str] = [] + drives: list = [] + try: + pool_map = await store.get_zfs_map() + drives = await get_host_drives(pool_map) + await store.set_host_drives(drives) + except Exception as e: + errors.append(f"hostdrives:{e}") + await store.set_meta({ + "last_sweep_ts": store.now(), + "last_sweep_duration": round(time.monotonic() - t0, 1), + "host_drive_count": len(drives), + "errors": errors[:20], + "interval": DRIVE_INTERVAL, + }, key=store.K_HOST_DRIVE_META) + logger.info("drive sweep: %d drives, %.1fs, %d error(s)", + len(drives), round(time.monotonic() - t0, 1), len(errors)) + + +async def run_loop(name: str, fn, interval: int, meta_key: str) -> None: + """Run `fn` every `interval`s, with a restart-storm guard from its meta.""" + meta = await store.get_meta(key=meta_key) + if meta and meta.get("last_sweep_ts"): + since = store.now() - meta["last_sweep_ts"] + if 0 <= since < interval: + wait = interval - since + logger.info("%s: last sweep %.0fs ago — deferring first %.0fs", name, since, wait) + await asyncio.sleep(wait) + while True: + t0 = time.monotonic() + try: + await fn() + except asyncio.CancelledError: + raise + except Exception as e: + logger.error("%s loop error: %s", name, e) + await asyncio.sleep(max(5, interval - (time.monotonic() - t0))) async def lock_refresher() -> None: @@ -117,30 +147,17 @@ async def main() -> None: if os.geteuid() != 0: logger.warning("not running as root — ipmitool/smartctl may fail") - jitter = random.uniform(0, STARTUP_JITTER) - logger.info("startup jitter %.1fs", jitter) - await asyncio.sleep(jitter) - - # Restart-storm guard (mirrors the SAS poller). - meta = await store.get_meta(key=store.K_HOST_META) - if meta and meta.get("last_sweep_ts"): - since = store.now() - meta["last_sweep_ts"] - if 0 <= since < SWEEP_INTERVAL: - wait = SWEEP_INTERVAL - since - logger.info("last host sweep %.0fs ago — deferring first sweep %.0fs", since, wait) - await asyncio.sleep(wait) + await asyncio.sleep(random.uniform(0, STARTUP_JITTER)) + env_task = asyncio.create_task( + run_loop("env", sweep_env, ENV_INTERVAL, store.K_HOST_ENV_META)) + drive_task = asyncio.create_task( + run_loop("drive", sweep_drives, DRIVE_INTERVAL, store.K_HOST_DRIVE_META)) try: - while True: - t0 = time.monotonic() - try: - await sweep() - except Exception as e: - logger.error("host sweep error: %s", e) - elapsed = time.monotonic() - t0 - await asyncio.sleep(max(5, SWEEP_INTERVAL - elapsed)) + await asyncio.gather(env_task, drive_task) finally: - refresher.cancel() + for t in (refresher, env_task, drive_task): + t.cancel() await close_cache() diff --git a/services/store.py b/services/store.py index 6ed0e1a..f3ec0a9 100644 --- a/services/store.py +++ b/services/store.py @@ -25,7 +25,8 @@ K_INVENTORY = "jbod:inventory" K_HOST_DRIVES = "jbod:host_drives" K_HOST_SENSORS = "jbod:host_sensors" K_META = "jbod:poll:meta" -K_HOST_META = "jbod:host_poll:meta" +K_HOST_ENV_META = "jbod:host_poll:env_meta" +K_HOST_DRIVE_META = "jbod:host_poll:drive_meta" K_LED_QUEUE = "jbod:led:queue" K_LOCK = "jbod:poll:lock" K_HOST_LOCK = "jbod:host_poll:lock" -- 2.49.1 From 514502e045071e080da95ebf205dac641b246352 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 19:05:10 +0000 Subject: [PATCH 09/13] host_sensors: filter IPMI margin sensors; coretemp per-core CPU fallback - skip IPMI '*Margin' sensors (thermal margin / distance-to-throttle reads in degrees C but is not an absolute temperature, e.g. robin's P1 Therm Margin) - coretemp: when a CPU has no 'Package id' sensor (older Xeons like robin's X3430), fall back to the hottest Core as the CPU temp --- services/host_sensors.py | 28 +++++++++++++++++++--------- 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/services/host_sensors.py b/services/host_sensors.py index 9bdc191..4736f57 100644 --- a/services/host_sensors.py +++ b/services/host_sensors.py @@ -45,6 +45,10 @@ def _parse_ipmi(text: str) -> list[dict]: if len(parts) < 5: continue name, _sid, status, entity, reading = parts[:5] + # Skip relative/derived sensors (e.g. "P1 Therm Margin") — they read in + # "degrees C" but are a distance-to-throttle, not an absolute temp. + if "margin" in name.lower(): + continue m = re.search(r"(-?\d+(?:\.\d+)?)\s*degrees?\s*C", reading, re.IGNORECASE) if not m: # "no reading", "disabled", "ns", etc. continue @@ -74,22 +78,28 @@ def read_coretemp() -> list[dict]: continue except OSError: continue + packages: list[dict] = [] + cores: list[int] = [] for label_f in sorted(hw.glob("temp*_label")): try: label = label_f.read_text().strip() except OSError: continue - if not label.lower().startswith("package id"): - continue input_f = label_f.with_name(label_f.name.replace("_label", "_input")) try: - milli = int(input_f.read_text().strip()) + temp = round(int(input_f.read_text().strip()) / 1000) except (OSError, ValueError): continue - pkg = label.split("id")[-1].strip() - cpus.append({ - "name": f"CPU {pkg}", - "temperature_c": round(milli / 1000), - "kind": "cpu", - }) + low = label.lower() + if low.startswith("package id"): + pkg = label.split("id")[-1].strip() + packages.append({"name": f"CPU {pkg}", "temperature_c": temp, "kind": "cpu"}) + elif low.startswith("core"): + cores.append(temp) + # Prefer the per-package sensor; older CPUs (no package sensor) fall back + # to the hottest core as the CPU temp. + if packages: + cpus.extend(packages) + elif cores: + cpus.append({"name": "CPU", "temperature_c": max(cores), "kind": "cpu"}) return cpus -- 2.49.1 From 4c68bff9648ce95f2ccde56681560bd4861269d8 Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 19:41:51 +0000 Subject: [PATCH 10/13] Add k8s DaemonSet design for RKE2 nodes (Model A, per-node self-contained) Design artifact (NOT applied). Per-node pod = host-poller + redis + web sharing pod localhost; only host-poller privileged with hostPath /dev+/sys. No SAS poller (no enclosures). MQTT_NODE_ID from spec.nodeName; MQTT creds via VSO VaultStaticSecret reading secret/home_assistant (mqtt_user/pass -> MQTT_USERNAME/ PASSWORD). nodeSelector jbod-monitor=enabled to pin to pascal/currie/fermi. --- k8s/README.md | 62 +++++++++++++++++++ k8s/daemonset.yaml | 104 ++++++++++++++++++++++++++++++++ k8s/mqtt-vaultstaticsecret.yaml | 30 +++++++++ 3 files changed, 196 insertions(+) create mode 100644 k8s/README.md create mode 100644 k8s/daemonset.yaml create mode 100644 k8s/mqtt-vaultstaticsecret.yaml diff --git a/k8s/README.md b/k8s/README.md new file mode 100644 index 0000000..743af39 --- /dev/null +++ b/k8s/README.md @@ -0,0 +1,62 @@ +# jbod-monitor on RKE2 (DaemonSet) + +Per-node self-contained deployment for Kubernetes worker nodes that can't run +the docker-compose stack. Each opted-in node runs one pod (`host-poller` + +`redis` + `web`) and publishes its own Home Assistant device +(`JBOD Monitor ()`) over MQTT. + +This mirrors the docker-compose stack with **no code changes** — each node is +independent, exactly like the standalone hosts. + +## What it collects + +- **host-poller** (privileged, hostPath `/dev`+`/sys`): IPMI/iDRAC env, CPU + coretemp, on-metal drive SMART → node-local redis. +- **web**: reads redis, publishes per-node MQTT discovery to the broker. +- No SAS poller (no SES enclosures on these nodes). Consequence: host-drive + entities have no ZFS pool label (no zfs_map producer). + +## Prereqs to verify before applying + +1. **Registry pull**: confirm the cluster can pull `docker.adamksmith.xyz/*`. + If it needs auth, sync a `regcred` into the `jbod-monitor` namespace and + uncomment `imagePullSecrets` in `daemonset.yaml`. +2. **VSO**: `vault-secrets-operator-system/openbao-kubernetes` VaultAuth exists + (it does, per infra notes). Confirm the `transformation.templates` syntax + matches the installed VSO version. +3. **MQTT egress**: pods must reach `10.5.30.3:1883` (EMQX). + +## Apply order + +```bash +# 1. Label ONLY the intended nodes (do not blanket-label all workers): +kubectl label node usco-dc-pascal usco-dc-currie usco-dc-fermi jbod-monitor=enabled + +# 2. Namespace + DaemonSet: +kubectl apply -f k8s/daemonset.yaml + +# 3. MQTT creds (VSO → Secret jbod-mqtt). Apply, then confirm it synced: +kubectl apply -f k8s/mqtt-vaultstaticsecret.yaml +kubectl -n jbod-monitor get vaultstaticsecret jbod-mqtt # SYNCED=True +kubectl -n jbod-monitor get secret jbod-mqtt # has MQTT_USERNAME/PASSWORD + +# 4. Roll the web container so it picks up the secret (if it started first): +kubectl -n jbod-monitor rollout restart daemonset/jbod-monitor +``` + +## Verify + +```bash +kubectl -n jbod-monitor get pods -o wide # one per labeled node +kubectl -n jbod-monitor logs -c host-poller | grep sweep +kubectl -n jbod-monitor logs -c web | grep -i mqtt # "MQTT connected" +``` + +Then check Home Assistant for `JBOD Monitor (usco-dc-pascal|currie|fermi)`. + +## Notes + +- pascal (R620) already has an iDRAC "System Board Inlet Temp" in HA via another + integration — its env temp will partially duplicate; drives + CPU are net-new. +- Pin image tags by SHA (already done: host-poller `514502e`, web `1839e9d`). +- Removal: `kubectl delete -f k8s/daemonset.yaml` and unlabel the nodes. diff --git a/k8s/daemonset.yaml b/k8s/daemonset.yaml new file mode 100644 index 0000000..8f47d84 --- /dev/null +++ b/k8s/daemonset.yaml @@ -0,0 +1,104 @@ +# jbod-monitor on RKE2 nodes — per-node self-contained DaemonSet (Model A). +# +# Each scheduled node runs one pod = host-poller + redis + web, sharing the +# pod's localhost (no hostNetwork needed). Only host-poller is privileged and +# mounts /dev + /sys to read IPMI / SMART / coretemp. web egresses to the MQTT +# broker and publishes a per-node HA device (MQTT_NODE_ID = spec.nodeName). +# +# No SAS poller (these nodes have no SES enclosures). NOTE: without it there is +# no zfs_map, so host-drive entities won't carry a ZFS pool label. +# +# Target nodes are chosen by label — see k8s/README.md (label only +# pascal/currie/fermi). Apply order: namespace -> VaultStaticSecret -> this. +--- +apiVersion: v1 +kind: Namespace +metadata: + name: jbod-monitor +--- +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: jbod-monitor + namespace: jbod-monitor + labels: + app: jbod-monitor +spec: + selector: + matchLabels: + app: jbod-monitor + template: + metadata: + labels: + app: jbod-monitor + spec: + # Only nodes explicitly opted in (see README: kubectl label node ...). + nodeSelector: + jbod-monitor: "enabled" + # If the registry needs auth, sync a regcred and uncomment: + # imagePullSecrets: + # - name: regcred + containers: + - name: redis + image: redis:7-alpine + args: ["redis-server", "--save", "60", "1", "--loglevel", "warning"] + volumeMounts: + - name: redis-data + mountPath: /data + resources: + requests: {cpu: 25m, memory: 32Mi} + limits: {memory: 128Mi} + + - name: host-poller + image: docker.adamksmith.xyz/jbod-host-poller:514502e + securityContext: + privileged: true # raw /dev access for ipmitool + smartctl + env: + - {name: TZ, value: "America/Denver"} + - {name: REDIS_HOST, value: "127.0.0.1"} + - {name: REDIS_PORT, value: "6379"} + - {name: HOST_ENV_INTERVAL, value: "60"} + - {name: HOST_DRIVE_INTERVAL, value: "300"} + - {name: POLL_CONCURRENCY, value: "1"} + - {name: POLL_DRIVE_GAP, value: "0.5"} + - {name: POLL_STARTUP_JITTER, value: "10"} + - {name: POLL_DATA_TTL, value: "900"} + volumeMounts: + - {name: dev, mountPath: /dev} + - {name: sys, mountPath: /sys, readOnly: true} + - {name: udev, mountPath: /run/udev, readOnly: true} + resources: + requests: {cpu: 25m, memory: 64Mi} + limits: {memory: 256Mi} + + - name: web + image: docker.adamksmith.xyz/jbod-web:1839e9d + env: + - {name: TZ, value: "America/Denver"} + - {name: REDIS_HOST, value: "127.0.0.1"} + - {name: REDIS_PORT, value: "6379"} + - {name: MQTT_HOST, value: "10.5.30.3"} + - {name: MQTT_PORT, value: "1883"} + - {name: MQTT_DISCOVERY_PREFIX, value: "homeassistant"} + - {name: MQTT_BASE_TOPIC, value: "jbod-monitor"} + - {name: MQTT_PUBLISH_INTERVAL, value: "60"} + - name: MQTT_NODE_ID + valueFrom: + fieldRef: + fieldPath: spec.nodeName + envFrom: + - secretRef: + name: jbod-mqtt # MQTT_USERNAME / MQTT_PASSWORD (see VaultStaticSecret) + resources: + requests: {cpu: 25m, memory: 64Mi} + limits: {memory: 256Mi} + + volumes: + - name: redis-data + emptyDir: {} + - name: dev + hostPath: {path: /dev} + - name: sys + hostPath: {path: /sys} + - name: udev + hostPath: {path: /run/udev} diff --git a/k8s/mqtt-vaultstaticsecret.yaml b/k8s/mqtt-vaultstaticsecret.yaml new file mode 100644 index 0000000..b1ac7de --- /dev/null +++ b/k8s/mqtt-vaultstaticsecret.yaml @@ -0,0 +1,30 @@ +# MQTT broker creds for the web container, synced from OpenBao by VSO. +# +# Reads secret/home_assistant (keys mqtt_user / mqtt_pass) and renders a k8s +# Secret `jbod-mqtt` with MQTT_USERNAME / MQTT_PASSWORD keys (consumed via +# envFrom in the DaemonSet). VSO's `vso-read` policy can read secret/data/*, +# so unlike the claude-read token it CAN read home_assistant — no creds ever +# touch these manifests. +# +# Verify the transformation syntax against the installed VSO version. +apiVersion: secrets.hashicorp.com/v1beta1 +kind: VaultStaticSecret +metadata: + name: jbod-mqtt + namespace: jbod-monitor +spec: + vaultAuthRef: vault-secrets-operator-system/openbao-kubernetes + mount: secret + type: kv-v2 + path: home_assistant + refreshAfter: 1h + destination: + name: jbod-mqtt + create: true + transformation: + excludeRaw: true + templates: + MQTT_USERNAME: + text: '{{ .Secrets.mqtt_user }}' + MQTT_PASSWORD: + text: '{{ .Secrets.mqtt_pass }}' -- 2.49.1 From 5032c58f7df804502a8f1e04aafc996ca6ea249c Mon Sep 17 00:00:00 2001 From: adam Date: Tue, 16 Jun 2026 19:46:39 +0000 Subject: [PATCH 11/13] host_sensors: read extra hwmon chips (dell_smm) as env sensors Boxes without IPMI (Dell workstations like dev-linux) expose CPU/ambient/ SODIMM/board temps via the BIOS SMM interface (dell_smm hwmon). Add read_hwmon_env(): labelled temps from chips in HOST_HWMON_CHIPS (default 'dell_smm') as env sensors, named '