1st POC
This commit is contained in:
@@ -0,0 +1,51 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
|
||||
from .data import extract_samples
|
||||
|
||||
LOG = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Collector:
|
||||
def __init__(self, client, store, module_names: dict[str, str], interval: int = 600):
|
||||
self.client = client
|
||||
self.store = store
|
||||
self.module_names = module_names
|
||||
self.interval = interval
|
||||
self._stop = threading.Event()
|
||||
self._thread: threading.Thread | None = None
|
||||
|
||||
def collect_once(self) -> int:
|
||||
samples = extract_samples(self.client.stations_data(), self.module_names)
|
||||
for name, sample in samples.items():
|
||||
try:
|
||||
self.store.update(name, sample)
|
||||
except Exception:
|
||||
# One stale/duplicate module must not discard the other modules.
|
||||
LOG.exception("Could not update %s RRD", name)
|
||||
LOG.info("Collected %d Netatmo modules", len(samples))
|
||||
return len(samples)
|
||||
|
||||
def start(self) -> None:
|
||||
if self._thread and self._thread.is_alive():
|
||||
return
|
||||
self._thread = threading.Thread(target=self._loop, name="netatmo-collector", daemon=True)
|
||||
self._thread.start()
|
||||
|
||||
def stop(self) -> None:
|
||||
self._stop.set()
|
||||
|
||||
def _loop(self) -> None:
|
||||
# Collect immediately, then align roughly to the configured cadence.
|
||||
while not self._stop.is_set():
|
||||
started = time.monotonic()
|
||||
try:
|
||||
self.collect_once()
|
||||
except Exception:
|
||||
LOG.exception("Netatmo collection failed")
|
||||
remaining = max(1, self.interval - (time.monotonic() - started))
|
||||
self._stop.wait(remaining)
|
||||
|
||||
Reference in New Issue
Block a user