"""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_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) ───────────────────────────────────────────────────── 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_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(key, 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_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 ──────────────────────────────────────────── 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 or the push fails (so the API can surface 503).""" client = get_client() if client is None: return False 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 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 (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, 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, key, token, ttl)) 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, key, token, ttl)) def now() -> float: return time.time()