Files
jbod-monitor/services/enclosure.py
adam 0f829e0380 Split into dedicated hardware poller + read-only consumers
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
2026-06-16 15:12:34 +00:00

348 lines
12 KiB
Python

import asyncio
import logging
import os
import re
from pathlib import Path
from services.hwgate import gate
logger = logging.getLogger(__name__)
ENCLOSURE_BASE = Path("/sys/class/enclosure")
def _read_sysfs(path: Path) -> str:
"""Read a sysfs attribute file, return stripped content or empty string."""
try:
return path.read_text().strip()
except (OSError, IOError):
return ""
def _find_sg_device(enclosure_path: Path) -> str | None:
"""Resolve the sg device for an enclosure from its sysfs path."""
# The enclosure sysfs directory has a 'device' symlink. Under that,
# there's a scsi_generic directory containing the sg device name.
sg_dir = enclosure_path / "device" / "scsi_generic"
if sg_dir.is_dir():
entries = list(sg_dir.iterdir())
if entries:
return f"/dev/{entries[0].name}"
return None
def discover_enclosures() -> list[dict]:
"""Walk /sys/class/enclosure/ to discover SES enclosures."""
if not ENCLOSURE_BASE.is_dir():
logger.warning("No enclosure sysfs directory found at %s", ENCLOSURE_BASE)
return []
enclosures = []
for enc_dir in sorted(ENCLOSURE_BASE.iterdir()):
if not enc_dir.is_dir():
continue
enc_id = enc_dir.name
device_dir = enc_dir / "device"
vendor = _read_sysfs(device_dir / "vendor")
model = _read_sysfs(device_dir / "model")
revision = _read_sysfs(device_dir / "rev")
sg_device = _find_sg_device(enc_dir)
slots = list_slots(enc_id)
total = len(slots)
populated = sum(1 for s in slots if s["populated"])
has_devices = any(s["device"] for s in slots)
# Skip secondary/passive IOM paths: enclosures where SES reports
# populated slots but the kernel has no device symlinks for any of
# them (the drives are owned by the primary IOM enclosure entry).
if populated > 0 and not has_devices:
logger.info(
"Skipping enclosure %s (%s %s) — %d populated slots but no "
"device links (likely a secondary IOM path)",
enc_id, vendor, model, populated,
)
continue
enclosures.append({
"id": enc_id,
"sg_device": sg_device,
"vendor": vendor,
"model": model,
"revision": revision,
"total_slots": total,
"populated_slots": populated,
})
return enclosures
def list_slots(enclosure_id: str) -> list[dict]:
"""Enumerate drive slots for an enclosure via sysfs."""
enc_dir = ENCLOSURE_BASE / enclosure_id
if not enc_dir.is_dir():
return []
slots = []
for entry in sorted(enc_dir.iterdir()):
if not entry.is_dir():
continue
# Determine if this is a drive slot element.
# Some enclosures use named dirs ("Slot 00", "Disk 1", "ArrayDevice00"),
# others use bare numeric dirs ("0", "1", "2") with a "type" file.
slot_num = _parse_slot_number(entry)
if slot_num is None:
continue
# Check if a block device is linked in this slot
block_dir = entry / "device" / "block"
device = None
populated = False
if block_dir.is_dir():
devs = list(block_dir.iterdir())
if devs:
device = devs[0].name
populated = True
else:
# Also check the 'status' file — "not installed" means empty
status = _read_sysfs(entry / "status")
if status and status not in ("not installed", ""):
populated = True
slots.append({
"slot": slot_num,
"populated": populated,
"device": device,
})
slots.sort(key=lambda s: s["slot"])
return slots
def _parse_slot_number(entry: Path) -> int | None:
"""Extract the slot number from a sysfs slot directory.
Handles multiple naming conventions:
- Bare numeric dirs ("0", "1") with type=device and a slot file
- Named dirs ("Slot 00", "Slot00", "Disk 1", "ArrayDevice00")
"""
name = entry.name
# Bare numeric directory — check the type file to confirm it's a device slot
if name.isdigit():
entry_type = _read_sysfs(entry / "type")
if entry_type not in ("device", "disk", "array device"):
return None
# Prefer the 'slot' file for the actual slot number
slot_val = _read_sysfs(entry / "slot")
if slot_val.isdigit():
return int(slot_val)
return int(name)
# Named directory prefixes
for prefix in ("Slot ", "Slot", "Disk ", "Disk", "ArrayDevice", "SLOT "):
if name.startswith(prefix):
num_str = name[len(prefix):].strip()
try:
return int(num_str)
except ValueError:
return None
return 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.
Hardware-gated."""
try:
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",
page, sg_device, stderr.decode().strip(),
)
return None
return stdout.decode(errors="replace")
except FileNotFoundError:
logger.warning("sg_ses not found")
return None
except Exception as e:
logger.warning("sg_ses page %s error for %s: %s", page, sg_device, e)
return None
async def get_enclosure_status(sg_device: str) -> dict | None:
"""Run sg_ses status page (0x02) and parse enclosure health data.
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 = await _run_sg_ses(sg_device, "0x02")
if status_text is None:
return None
return _parse_ses_page02(status_text)
def _norm_type(element_type: str) -> str:
"""Normalize an SES element type line to a stable key.
sg_ses appends qualifiers like ", subenclosure id: 0 [ti=2]" to the
type line; strip those so the same element type matches across the
status (0x02) and element descriptor (0x07) pages.
"""
return element_type.split(",")[0].strip().lower()
def _parse_element_descriptors(text: str) -> dict[tuple[str, int], str]:
"""Parse sg_ses --page=0x07 into {(element_type, index): name}."""
names: dict[tuple[str, int], str] = {}
sections = re.split(r"(?=\s*Element type:)", text)
for section in sections:
type_match = re.match(r"\s*Element type:\s*(.+)", section)
if not type_match:
continue
etype = _norm_type(type_match.group(1).strip().rstrip(","))
for line in section.splitlines():
m = re.match(r"\s*Element (\d+) descriptor:\s*(.*)", line)
if not m:
continue
idx = int(m.group(1))
name = m.group(2).strip()
# sg_ses prints "<empty>" when an enclosure leaves the descriptor
# string blank — treat that as no name so callers fall back to a
# generic label.
if name and name.lower() != "<empty>":
names[(etype, idx)] = name
return names
def _parse_ses_page02(
text: str, names: dict[tuple[str, int], str] | None = None
) -> dict:
"""Parse sg_ses --page=0x02 text output into structured health data.
``names`` maps (normalized element type, index) -> descriptor name from
the element descriptor page (0x07); when present it labels each element.
"""
names = names or {}
result = {
"overall_status": "OK",
"psus": [],
"fans": [],
"temps": [],
"voltages": [],
}
# Split into element type sections.
# Each section starts with "Element type: <type>"
sections = re.split(r"(?=\s*Element type:)", text)
for section in sections:
type_match = re.match(r"\s*Element type:\s*(.+)", section)
if not type_match:
continue
element_type = type_match.group(1).strip().rstrip(",").lower()
etype_key = _norm_type(element_type)
# Find individual element blocks (skip "Overall descriptor")
elements = re.split(r"(?=\s*Element \d+ descriptor:)", section)
for elem_text in elements:
desc_match = re.match(r"\s*Element (\d+) descriptor:", elem_text)
if not desc_match:
continue
idx = int(desc_match.group(1))
# Extract status line
status_match = re.search(r"status:\s*(.+?)(?:,|\n|$)", elem_text, re.IGNORECASE)
status = status_match.group(1).strip() if status_match else "Unknown"
if status.lower() == "not installed":
continue
name = names.get((etype_key, idx))
if "power supply" in element_type:
fail = "Fail=1" in elem_text
ac_fail = "AC fail=1" in elem_text
dc_fail = "DC fail=1" in elem_text
result["psus"].append({
"index": idx,
"name": name,
"status": status,
"fail": fail,
"ac_fail": ac_fail,
"dc_fail": dc_fail,
})
elif "cooling" in element_type or "fan" in element_type:
fail = "Fail=1" in elem_text
rpm_match = re.search(r"Actual speed[=:]\s*(\d+)\s*rpm", elem_text, re.IGNORECASE)
rpm = int(rpm_match.group(1)) if rpm_match else None
result["fans"].append({
"index": idx,
"name": name,
"status": status,
"rpm": rpm,
"fail": fail,
})
elif "temperature" in element_type:
temp_match = re.search(r"Temperature=\s*([\d.]+)\s*C", elem_text)
temp = float(temp_match.group(1)) if temp_match else None
# Skip dead/disconnected sensors: 0°C with a non-OK status
# is a non-functional sensor slot, not an actual reading.
if (temp is None or temp == 0) and status.lower() in (
"unrecoverable", "unknown", "not available",
):
continue
result["temps"].append({
"index": idx,
"name": name,
"status": status,
"temperature_c": temp,
})
elif "voltage" in element_type:
volt_match = re.search(r"Voltage:\s*([\d.]+)\s*V", elem_text, re.IGNORECASE)
if not volt_match:
volt_match = re.search(r"([\d.]+)\s*V", elem_text)
voltage = float(volt_match.group(1)) if volt_match else None
result["voltages"].append({
"index": idx,
"name": name,
"status": status,
"voltage": voltage,
})
# Derive overall_status from the actual parsed elements rather than
# the raw SES header, which counts dead/disconnected sensors as
# UNRECOV and inflates severity.
all_statuses = (
[e["status"] for e in result["psus"]]
+ [e["status"] for e in result["fans"]]
+ [e["status"] for e in result["temps"]]
+ [e["status"] for e in result["voltages"]]
)
status_lower = [s.lower() for s in all_statuses]
if any(s in ("unrecoverable", "critical") for s in status_lower):
result["overall_status"] = "CRITICAL"
elif any(s in ("noncritical", "non-critical", "warning") for s in status_lower):
result["overall_status"] = "WARNING"
elif any(s not in ("ok", "unknown") for s in status_lower):
result["overall_status"] = "WARNING"
return result