From a88600523a52cc8570422913055635f8137356ad Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Wed, 16 Sep 2026 18:29:27 -0700 Subject: [PATCH] Add the keygrabber control keywords and a testable scheduler --- docs/source/api/index.md | 6 + docs/source/keygrabber.md | 57 ++++++- libby/keygrabber/collection.py | 15 +- libby/keygrabber/daemon.py | 256 +++++++++++++++++++++++------ libby/keygrabber/scheduler.py | 114 +++++++++++++ tests/test_keygrabber_daemon.py | 109 +++++++++++- tests/test_keygrabber_scheduler.py | 169 +++++++++++++++++++ 7 files changed, 668 insertions(+), 58 deletions(-) create mode 100644 libby/keygrabber/scheduler.py create mode 100644 tests/test_keygrabber_scheduler.py diff --git a/docs/source/api/index.md b/docs/source/api/index.md index d6f2389..55bbb98 100644 --- a/docs/source/api/index.md +++ b/docs/source/api/index.md @@ -85,6 +85,12 @@ the docs build does not install the optional `influxdb` extra. :show-inheritance: ``` +```{eval-rst} +.. automodule:: libby.keygrabber.scheduler + :members: + :show-inheritance: +``` + ```{eval-rst} .. automodule:: libby.keygrabber.daemon :members: diff --git a/docs/source/keygrabber.md b/docs/source/keygrabber.md index e75a85e..da5968f 100644 --- a/docs/source/keygrabber.md +++ b/docs/source/keygrabber.md @@ -20,10 +20,9 @@ pip install libby[influxdb] keygrabber -c /etc/hispec/keygrabber.yaml ``` -It is an ordinary `LibbyDaemon`, so `SIGTERM` stops it cleanly and the -`shutdown` keyword will too once the control surface lands. On the way out it -gives its retry queue a bounded chance to drain, so a graceful stop does not -lose the last tick. +It is an ordinary `LibbyDaemon`, so `SIGTERM` stops it cleanly and so does +writing its `shutdown` keyword. On the way out it gives its retry queue a +bounded chance to drain, so a graceful stop does not lose the last tick. ## Config @@ -99,6 +98,56 @@ and a tighter interval would be overrun by a single slow peer. A tick whose predecessor is still running is skipped rather than queued behind it, so a wedged peer cannot accumulate overlapping reads. +## Control keywords + +The keygrabber is itself a peer, so its cadence and health are reachable with +the ordinary `libby` verbs, and it can feed the same dashboards it fills. + +| Keyword | Type | Access | Meaning | +|---|---|---|---| +| `enabled` | bool | R/W | Collect on the configured cadences; false pauses without exiting | +| `isconnected` | bool | R/W | Last sink write succeeded; write true to request a reconnect | +| `pointswritten` | int | R | Samples the sink has stored since start | +| `readerrors` | int | R | Keyword reads that failed | +| `writeerrors` | int | R | Sink writes that failed | +| `queuedepth` | int | R | Batches waiting in the retry queue | +| `skippedticks` | int | R | Ticks skipped because the previous read was still in flight | +| `droppedbatches` | int | R | Batches discarded because a queue was full | +| `reload` | trigger | W | Re-read the config file and apply it | +| `shutdown` | trigger | W | Gracefully stop the daemon | +| `.enabled` | bool | R/W | Collect this one collection | +| `.interval` | float | R/W | Cadence in seconds | +| `.lastsample` | string | R | UTC time of the last successful tick | +| `.lag` | float | R | Seconds the last tick ran past its due time | + +```bash +libby show hispec.keygrabber.% +libby modify hispec.keygrabber.adc.interval=30 +libby modify hispec.keygrabber.enabled=false +``` + +Per-collection keywords contain a dot, and `%` matches within a single segment +only, so `libby list hispec.keygrabber.%` will not show `adc.enabled`. Use +`hispec.keygrabber.%.%` for those. + +`isconnected` reports whether the last write to the sink succeeded; it does not +ping the database, because these getters are answered on the transport's +receive thread and a blocking one would time out every read in flight. Writing +`true` asks the writer thread to reconnect and returns immediately, so poll the +keyword for the outcome. There is no manual disconnect. + +### reload + +`reload` re-reads the config file and applies it to the running collections, +so keyword selections and cadences can change without a restart. A file that +fails to parse leaves the running collections untouched and reports why, both +to the caller and on `lasterror`. + +It will not add or remove collections, and says so rather than half-applying: +libby has no way to withdraw a keyword, so a new collection's control keywords +could not appear without a restart. Changing `sink` or `workers` also needs a +restart. + ## How it reads A tick is one `keys.read` request per peer, not one per keyword. This matters diff --git a/libby/keygrabber/collection.py b/libby/keygrabber/collection.py index 59f744d..ea2e742 100644 --- a/libby/keygrabber/collection.py +++ b/libby/keygrabber/collection.py @@ -22,7 +22,9 @@ class TickResult: read_errors: int -class Collection: +# Config plus the runtime state the control keywords expose; each attribute is +# one reported value rather than hidden complexity +class Collection: # pylint: disable=too-many-instance-attributes """Tracks what one peer exposes and turns a read of it into samples. Resolution is refreshed periodically rather than once, so keywords added by @@ -36,6 +38,13 @@ def __init__( clock: Callable[[], float] = time.monotonic, ) -> None: self.config = config + # Runtime state the control keywords read and write. Plain attributes + # rather than lock-guarded: each is a single value written by one + # thread and read by the transport's receive thread, and that thread + # must never block on a lock a tick might hold. + self.enabled = True + self.last_sample: Optional[datetime] = None + self.lag_s = 0.0 self._clock = clock self._names: Tuple[str, ...] = () self._bulk_read = False @@ -62,6 +71,10 @@ def needs_resolve(self) -> bool: return True return self._clock() - self._resolved_at >= self.config.refresh_s + def invalidate(self) -> None: + """Force the next tick to resolve again, after a config change.""" + self._resolved_at = None + def resolve(self, client: Client) -> Tuple[str, ...]: """Ask the peer what it serves and select the configured keywords. diff --git a/libby/keygrabber/daemon.py b/libby/keygrabber/daemon.py index 4ad573c..a65c1c1 100644 --- a/libby/keygrabber/daemon.py +++ b/libby/keygrabber/daemon.py @@ -1,7 +1,7 @@ """The keygrabber daemon: schedule reads, collect samples, write them out.""" from __future__ import annotations -import heapq +import dataclasses import queue import threading import time @@ -10,11 +10,13 @@ from typing import Callable, List, Optional, Sequence, Tuple from ..client import Client +from ..config import ConfigError, DaemonConfigLoader, with_env_overrides from ..daemon import LibbyDaemon from ..errors import LibbyError from ..libby import Libby from .collection import Collection -from .config import KeygrabberConfig, build_sink, parse_config +from .config import TIMEOUT_HEADROOM, KeygrabberConfig, build_sink, parse_config +from .scheduler import Scheduler from .sink import RetryingWriter, Sample, Sink # How long the writer waits for a batch before checking the retry queue @@ -40,6 +42,15 @@ STOP_JOIN_TIMEOUT_S = DRAIN_DEADLINE_S + WRITER_POLL_S * 2 +@dataclasses.dataclass(frozen=True) +class _ConfigSource: + """Where a daemon's config came from, so ``reload`` can re-read it.""" + + path: str + daemon_id: Optional[str] = None + env_prefix: Optional[str] = None + + class Counters: # pylint: disable=too-few-public-methods """Tallies the daemon reports, guarded for cross-thread increments. @@ -79,6 +90,7 @@ class KeygrabberDaemon(LibbyDaemon): # pylint: disable=too-many-instance-attrib def __init__(self) -> None: super().__init__() self.counters = Counters() + self.enabled = True self._settings: Optional[KeygrabberConfig] = None self._client: Optional[Client] = None self._writer: Optional[RetryingWriter] = None @@ -88,8 +100,10 @@ def __init__(self) -> None: self._threads: List[threading.Thread] = [] self._writer_thread: Optional[threading.Thread] = None self._halt = threading.Event() - self._in_flight: set[str] = set() - self._in_flight_lock = threading.Lock() + self._scheduler = Scheduler() + self._source: Optional[_ConfigSource] = None + self._sink_healthy = True + self._reconnect_wanted = False self._clock: Callable[[], float] = time.monotonic def on_start(self, libby: Libby) -> None: @@ -105,6 +119,8 @@ def on_start(self, libby: Libby) -> None: maxsize=max(1, self._settings.workers * QUEUE_DEPTH_PER_WORKER)) self._pool = ThreadPoolExecutor(max_workers=self._settings.workers, thread_name_prefix="keygrabber-read") + self._scheduler.replace(self._collections) + self._register_control_keywords() self._halt.clear() self._writer_thread = self._spawn(self._write_loop, "keygrabber-write") self._spawn(self._schedule_loop, "keygrabber-schedule") @@ -137,6 +153,13 @@ def on_stop(self, libby: Optional[Libby] = None) -> None: return writer.close() + @classmethod + def from_config_file(cls, path, daemon_id=None, *, env_prefix=None): + """Build from a file, remembering where it came from for ``reload``.""" + daemon = super().from_config_file(path, daemon_id, env_prefix=env_prefix) + daemon._source = _ConfigSource(str(path), daemon_id, env_prefix) + return daemon + def make_sink(self) -> Sink: """Build the configured sink. @@ -147,6 +170,146 @@ def make_sink(self) -> Sink: raise LibbyError("keygrabber config has not been parsed yet") return build_sink(self._settings.sink) + ### Control surface + + def _register_control_keywords(self) -> None: + """Expose the daemon's own state as keywords. + + Every getter here is answered on the transport's receive thread, so + none of them may block: the counters are plain reads, ``isconnected`` + reports a flag the writer thread maintains rather than pinging the + database, and writing it only requests a reconnect. + """ + registry = self.keyword_registry + registry.bool("enabled", + getter=lambda: self.enabled, + setter=self._set_enabled, + description="Collect on the configured cadences; " + "write false to pause without exiting.") + registry.bool("isconnected", + getter=lambda: self._sink_healthy, + setter=self._set_connected, + description="Last sink write succeeded; write true to " + "request a reconnect.") + registry.int("pointswritten", getter=lambda: self.counters.points_written, + description="Samples the sink has stored since start.") + registry.int("readerrors", getter=lambda: self.counters.read_errors, + description="Keyword reads that failed.") + registry.int("writeerrors", getter=lambda: self.counters.write_errors, + description="Sink writes that failed.") + registry.int("queuedepth", getter=self._queue_depth, + description="Batches waiting in the retry queue.") + registry.int("skippedticks", getter=lambda: self.counters.skipped_ticks, + description="Ticks skipped because the previous read " + "was still in flight.") + registry.int("droppedbatches", getter=lambda: self.counters.dropped_batches, + description="Batches discarded because a queue was full.") + registry.trigger("reload", action=self._reload, + description="Re-read the config file and apply it.") + registry.trigger("shutdown", action=self.request_stop, + description="Gracefully stop this daemon.") + for collection in self._collections: + self._register_collection_keywords(collection) + + def _register_collection_keywords(self, collection: Collection) -> None: + """Expose one collection's cadence and health.""" + name = collection.name + registry = self.keyword_registry + registry.bool(f"{name}.enabled", + getter=lambda c=collection: c.enabled, + setter=lambda value, c=collection: setattr(c, "enabled", value), + description=f"Collect {name}; write false to pause it.") + registry.float(f"{name}.interval", + getter=lambda c=collection: c.config.interval_s, + setter=lambda value, n=name: self._set_interval(n, value), + units="seconds", + description=f"Cadence for {name}.") + registry.string(f"{name}.lastsample", + getter=lambda c=collection: ( + c.last_sample.isoformat() if c.last_sample else None), + nullable=True, + description=f"UTC time of the last successful {name} tick.") + registry.float(f"{name}.lag", + getter=lambda c=collection: c.lag_s, + units="seconds", + description=f"Seconds the last {name} tick ran past due.") + + def _set_enabled(self, value: bool) -> None: + """Pause or resume all collection.""" + self.enabled = value + self.logger.info("collection %s", "enabled" if value else "paused") + + def _set_connected(self, value: bool) -> None: + """Request a reconnect; refuse a write of false. + + The reconnect happens on the writer thread rather than here, because a + dead database would otherwise hold the receive thread for the client's + whole connect timeout. Poll ``isconnected`` for the outcome. + """ + if not value: + raise ValueError("write true to reconnect; there is no manual disconnect") + self._reconnect_wanted = True + + def _set_interval(self, name: str, value: float) -> None: + """Change one collection's cadence, keeping the timeout headroom rule.""" + collection = self._collection(name) + minimum = TIMEOUT_HEADROOM * collection.config.timeout_s + if value <= minimum: + raise ValueError( + f"interval must exceed {minimum}s, which is " + f"{TIMEOUT_HEADROOM} x this collection's timeout") + collection.config = dataclasses.replace(collection.config, + interval_s=value) + self._scheduler.update(collection) + + def _collection(self, name: str) -> Collection: + for collection in self._collections: + if collection.name == name: + return collection + raise ValueError(f"unknown collection {name!r}") + + def _queue_depth(self) -> int: + writer = self._writer + return writer.queue_depth if writer is not None else 0 + + def _reload(self) -> None: + """Re-read the config file and apply it to the running collections. + + Parsed here, on the receive thread, so a bad file is reported straight + back to the caller and the running set is left untouched. Adding or + removing a collection is refused rather than half-applied: libby has no + way to withdraw a keyword, so a new collection's control keywords could + not appear without a restart. + """ + if self._source is None: + raise LibbyError("this keygrabber was not built from a config file") + + config = DaemonConfigLoader(self._source.path).get_daemon_config( + self._source.daemon_id) + if self._source.env_prefix: + config = with_env_overrides(config, prefix=self._source.env_prefix) + settings = parse_config(config) + + running = {collection.name for collection in self._collections} + reloaded = {collection.name for collection in settings.collections} + if running != reloaded: + raise ConfigError( + "reload cannot add or remove collections " + f"({sorted(running)} -> {sorted(reloaded)}); restart instead" + ) + + by_name = {collection.name: collection for collection in self._collections} + for config_entry in settings.collections: + collection = by_name[config_entry.name] + collection.config = config_entry + # The keyword selection may have changed, so do not wait out the + # old refresh window before picking it up + collection.invalidate() + self._settings = settings + self._scheduler.replace(self._collections) + self.logger.info("reloaded %d collections from %s", + len(self._collections), self._source.path) + def _spawn(self, target: Callable[[], None], name: str) -> threading.Thread: thread = threading.Thread(target=target, name=name, daemon=True) thread.start() @@ -156,49 +319,26 @@ def _spawn(self, target: Callable[[], None], name: str) -> threading.Thread: ### Scheduling def _schedule_loop(self) -> None: - """Submit each collection's tick when it comes due.""" - now = self._clock() - pending: List[Tuple[float, str]] = [ - (now, collection.name) for collection in self._collections - ] - heapq.heapify(pending) - by_name = {collection.name: collection for collection in self._collections} - + """Submit each collection's tick when the scheduler says it is due.""" while not self._halt.is_set(): - if not pending: - self._halt.wait(SCHEDULER_TICK_S) - continue - - due_at, name = pending[0] - delay = due_at - self._clock() - if delay > 0: - self._halt.wait(min(delay, SCHEDULER_TICK_S)) - continue - - heapq.heappop(pending) - collection = by_name[name] - self._submit(collection) - # Never schedule into the past: a long stall would otherwise queue - # a burst of catch-up ticks that can only skip - interval = collection.config.interval_s - heapq.heappush(pending, (max(due_at + interval, self._clock()), name)) - - def _submit(self, collection: Collection) -> None: - """Run a tick unless the previous one is still going.""" - with self._in_flight_lock: - if collection.name in self._in_flight: - self.counters.add(skipped_ticks=1) - self.logger.warning( - "collection %s skipped: previous read still in flight", - collection.name) - return - self._in_flight.add(collection.name) - - if not self._dispatch(collection): - # Shut down between the claim and the submit, so give the claim - # back rather than leaving the collection marked busy for good - with self._in_flight_lock: - self._in_flight.discard(collection.name) + if self.enabled: + self._submit_due() + delay = self._scheduler.next_delay() + self._halt.wait(SCHEDULER_TICK_S if delay is None + else min(delay, SCHEDULER_TICK_S)) + + def _submit_due(self) -> None: + """Hand every due collection to the pool, counting the skips.""" + claimed, skipped = self._scheduler.claim_due() + if skipped: + self.counters.add(skipped_ticks=skipped) + self.logger.warning("%d tick(s) skipped: previous read still in flight", + skipped) + for collection in claimed: + if not collection.enabled or not self._dispatch(collection): + # Paused, or shut down between the claim and the submit: give + # the claim back rather than leaving it marked busy for good + self._scheduler.release(collection.name) def _dispatch(self, collection: Collection) -> bool: """Submit a tick, returning False once the pool can no longer take one.""" @@ -224,12 +364,12 @@ def _run_tick(self, collection: Collection) -> None: self.counters.add(read_errors=result.read_errors) if result.samples: self._enqueue(result.samples) + collection.last_sample = datetime.now(timezone.utc) except LibbyError as exc: self.counters.add(read_errors=1) self.logger.error("collection %s read failed: %s", collection.name, exc) finally: - with self._in_flight_lock: - self._in_flight.discard(collection.name) + self._scheduler.release(collection.name) def _enqueue(self, samples: Sequence[Sample]) -> None: try: @@ -244,6 +384,7 @@ def _enqueue(self, samples: Sequence[Sample]) -> None: def _write_loop(self) -> None: """Own every sink call, so no reader thread ever touches the backend.""" while not self._halt.is_set(): + self._reconnect_if_requested() try: batch = self._queue.get(timeout=WRITER_POLL_S) except queue.Empty: @@ -254,14 +395,32 @@ def _write_loop(self) -> None: # Halted: drain here, on the one thread allowed to touch the sink self._drain() + def _reconnect_if_requested(self) -> None: + """Honour an isconnected write, on this thread rather than the caller's.""" + if not self._reconnect_wanted: + return + self._reconnect_wanted = False + writer = self._writer + if writer is None: + return + try: + writer.connect() + self._sink_healthy = True + self.logger.info("reconnected to the sink on request") + except LibbyError as exc: + self._sink_healthy = False + self.logger.error("sink reconnect failed: %s", exc) + def _write(self, batch: Tuple[Sample, ...]) -> None: writer = self._writer if writer is None: return try: self.counters.add(points_written=writer.write(batch)) + self._sink_healthy = True except LibbyError as exc: self.counters.add(write_errors=1) + self._sink_healthy = False self.logger.error("sink write failed: %s", exc) def _flush(self) -> None: @@ -272,6 +431,7 @@ def _flush(self) -> None: self.counters.add(points_written=writer.flush_due()) except LibbyError as exc: self.counters.add(write_errors=1) + self._sink_healthy = False self.logger.error("sink retry failed: %s", exc) def _drain(self) -> None: diff --git a/libby/keygrabber/scheduler.py b/libby/keygrabber/scheduler.py new file mode 100644 index 0000000..f0ec5d4 --- /dev/null +++ b/libby/keygrabber/scheduler.py @@ -0,0 +1,114 @@ +"""Decides which collections are due to be read, and which to skip. + +Deliberately free of threads, pools and sleeping: it answers "what should run +now" and nothing else, so the due-time ordering and the skip rule can be tested +against an injected clock rather than by racing real threads. +""" +from __future__ import annotations + +import heapq +import threading +import time +from typing import Callable, Dict, Iterable, List, Optional, Set, Tuple + +from .collection import Collection + + +class Scheduler: + """A due-time heap over collections, with at most one tick each in flight. + + Every method takes the lock only for in-memory bookkeeping, never across a + read or a write, because ``reload`` calls :meth:`replace` from the + transport's receive thread and blocking that would time out every read + already in flight. + """ + + def __init__(self, *, clock: Callable[[], float] = time.monotonic) -> None: + self._clock = clock + self._lock = threading.Lock() + self._collections: Dict[str, Collection] = {} + self._due: List[Tuple[float, str]] = [] + self._in_flight: Set[str] = set() + + def replace(self, collections: Iterable[Collection]) -> None: + """Swap the scheduled set, leaving in-flight ticks to finish. + + Every new collection is due immediately, so a reload takes effect on + the next pass rather than after a full interval. + """ + now = self._clock() + with self._lock: + self._collections = {item.name: item for item in collections} + self._due = [(now, name) for name in self._collections] + heapq.heapify(self._due) + # Names no longer scheduled are dropped; a tick still running for + # one releases harmlessly below + self._in_flight &= set(self._collections) + + def next_delay(self) -> Optional[float]: + """Return seconds until the next tick, or None when nothing is scheduled.""" + with self._lock: + if not self._due: + return None + return max(0.0, self._due[0][0] - self._clock()) + + def claim_due(self) -> Tuple[Tuple[Collection, ...], int]: + """Claim every collection now due, and count those already running. + + Returns the collections the caller should tick, and how many ticks were + skipped because their predecessor had not finished. A skipped tick is + not retried sooner: it waits for its next interval, so a slow peer + settles at a lower rate instead of building a backlog. + """ + claimed: List[Collection] = [] + skipped = 0 + with self._lock: + now = self._clock() + while self._due and self._due[0][0] <= now: + due_at, name = heapq.heappop(self._due) + collection = self._collections.get(name) + if collection is None: + continue # dropped by a reload; stop scheduling it + if name in self._in_flight: + skipped += 1 + else: + self._in_flight.add(name) + collection.lag_s = max(0.0, now - due_at) + claimed.append(collection) + # Rescheduled from now, not from when it was due, so a long + # stall cannot leave a burst of catch-up ticks that can only + # skip. Cadence drifts by the loop's own latency instead. + heapq.heappush( + self._due, (now + collection.config.interval_s, name)) + return tuple(claimed), skipped + + def update(self, collection: Collection) -> None: + """Replace one collection and reschedule only it. + + Used when a control keyword changes a cadence, so adjusting one + collection does not re-seed every other collection's due time and set + the whole fleet reading at once. + """ + with self._lock: + self._collections[collection.name] = collection + self._due = [entry for entry in self._due + if entry[1] != collection.name] + heapq.heapify(self._due) + heapq.heappush(self._due, (self._clock(), collection.name)) + + def release(self, name: str) -> None: + """Mark a collection's tick finished, so its next one may run.""" + with self._lock: + self._in_flight.discard(name) + + @property + def in_flight(self) -> int: + """Number of ticks currently running.""" + with self._lock: + return len(self._in_flight) + + @property + def scheduled(self) -> Tuple[str, ...]: + """Names currently scheduled, in no particular order.""" + with self._lock: + return tuple(self._collections) diff --git a/tests/test_keygrabber_daemon.py b/tests/test_keygrabber_daemon.py index e2834a9..db33819 100644 --- a/tests/test_keygrabber_daemon.py +++ b/tests/test_keygrabber_daemon.py @@ -15,6 +15,7 @@ import unittest from typing import List, Sequence +from libby import Client, KeywordError from libby.daemon import LibbyDaemon from libby.keygrabber import KeygrabberDaemon, Sample from libby.rabbitmq_transport import RabbitMQTransport @@ -195,12 +196,98 @@ def test_submit_to_a_shut_down_pool_releases_the_claim(self): self.grabber.start() self._await_samples() # pylint: disable=protected-access - collection = self.grabber._collections[0] + scheduler = self.grabber._scheduler self.grabber._pool.shutdown(wait=True) # left non-None on purpose - self.grabber._in_flight.discard(collection.name) - self.grabber._submit(collection) # must not raise - self.assertNotIn(collection.name, self.grabber._in_flight) + self.grabber._submit_due() # must not raise + self.assertEqual(scheduler.in_flight, 0) + + def _grabber_client(self) -> Client: + """Return a client addressed at the keygrabber's own keywords.""" + raise NotImplementedError + + def _keyword(self, name: str) -> str: + return f"hispec.{self.grabber.peer_id}.{name}" + + def test_control_keywords_answer_while_collecting(self): + """Serve the daemon's own keywords while ticks are in flight. + + The getters run on the transport's receive thread, the same thread + that delivers replies to the reader threads, so one that blocked + would time out every read in flight. + """ + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + self.assertTrue(client.get(self._keyword("enabled"))) + self.assertTrue(client.get(self._keyword("isconnected"))) + self.assertGreater(client.get(self._keyword("pointswritten")), 0) + self.assertEqual(client.get(self._keyword("writeerrors")), 0) + self.assertEqual(client.get(self._keyword("queuedepth")), 0) + self.assertIsInstance(client.get(self._keyword("readerrors")), int) + + def test_pausing_stops_collection_without_exiting(self): + """Stop collecting on a write of false, and resume on true.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + client.set(self._keyword("enabled"), False) + self.assertFalse(client.get(self._keyword("enabled"))) + time.sleep(1.2) # more than the 0.5s cadence + paused_at = len(self.sink.keywords()) + time.sleep(1.2) + self.assertEqual(len(self.sink.keywords()), paused_at) + + client.set(self._keyword("enabled"), True) + deadline = time.monotonic() + SETTLE_TIMEOUT_S + while time.monotonic() < deadline: + if len(self.sink.keywords()) > paused_at: + return + time.sleep(0.05) + self.fail("collection did not resume after being re-enabled") + + def test_per_collection_cadence_is_adjustable(self): + """Change one collection's interval over the wire.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + self.assertEqual(client.get(self._keyword("target.interval")), 0.5) + self.assertEqual(client.set(self._keyword("target.interval"), 5.0), + 5.0) + self.assertEqual(client.get(self._keyword("target.interval")), 5.0) + + def test_cadence_below_the_timeout_headroom_is_refused(self): + """Refuse a live cadence change the config loader would reject.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + with self.assertRaises(KeywordError): + client.set(self._keyword("target.interval"), 0.3) + + def test_collection_health_is_reported(self): + """Report the last tick's time and lateness per collection.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + self.assertIsNotNone(client.get(self._keyword("target.lastsample"))) + self.assertGreaterEqual(client.get(self._keyword("target.lag")), 0.0) + + def test_disconnect_is_refused_but_reconnect_is_accepted(self): + """Offer a reconnect without offering a manual disconnect.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + with self.assertRaises(KeywordError): + client.set(self._keyword("isconnected"), False) + self.assertTrue(client.set(self._keyword("isconnected"), True)) + + def test_reload_without_a_config_file_is_refused(self): + """Report that a daemon built from a mapping has nothing to re-read.""" + self.grabber.start() + self._await_samples() + with self._grabber_client() as client: + with self.assertRaises(KeywordError): + client.set(self._keyword("reload"), 1) def test_repeats_on_the_configured_cadence(self): """Read again on the next interval rather than once at startup.""" @@ -240,6 +327,10 @@ class ZmqKeygrabberTests(_Bases.KeygrabberCases): target_peer = f"hsfei.{_ZmqFixtureDaemon.peer_id}" + def setUp(self) -> None: + self.grabber_endpoint = _free_endpoint() + super().setUp() + @classmethod def setUpClass(cls): cls.endpoint = _free_endpoint() @@ -251,12 +342,17 @@ def setUpClass(cls): def tearDownClass(cls): cls.target.stop() + def _grabber_client(self) -> Client: + return Client.zmq(bind=_free_endpoint(), + address_book={f"hispec.{self.grabber.peer_id}": + self.grabber_endpoint}) + def grabber_config(self) -> dict: return { "peer_id": "keygrabberzmq", "group_id": "hispec", "transport": "zmq", - "bind": _free_endpoint(), + "bind": self.grabber_endpoint, "address_book": {self.target_peer: self.endpoint}, "discovery_enabled": False, "sink": {"type": "recording"}, @@ -280,6 +376,9 @@ def setUpClass(cls): def tearDownClass(cls): cls.target.stop() + def _grabber_client(self) -> Client: + return Client.rabbitmq(rabbitmq_url=RABBITMQ_URL) + def grabber_config(self) -> dict: return { "peer_id": "keygrabberrmq", diff --git a/tests/test_keygrabber_scheduler.py b/tests/test_keygrabber_scheduler.py new file mode 100644 index 0000000..b8156e2 --- /dev/null +++ b/tests/test_keygrabber_scheduler.py @@ -0,0 +1,169 @@ +"""Unit tests for Scheduler: due ordering, the skip rule and rescheduling. + +The clock is injected, so none of this waits for a real interval. +""" +from __future__ import annotations + +import dataclasses +import unittest + +from libby.keygrabber import Collection +from libby.keygrabber.config import CollectionConfig +from libby.keygrabber.scheduler import Scheduler + + +def _collection(name: str, interval_s: float = 10.0) -> Collection: + return Collection(CollectionConfig( + name=name, group="hsfei", daemon=name, keywords=("%",), exclude=(), + interval_s=interval_s, timeout_s=1.0, refresh_s=300.0)) + + +class _FakeClock: + """Manually advanced clock.""" + + def __init__(self) -> None: + self.now = 0.0 + + def __call__(self) -> float: + return self.now + + def advance(self, seconds: float) -> None: + """Move the clock forward.""" + self.now += seconds + + +class SchedulerTests(unittest.TestCase): + """What the scheduler decides to run, and when.""" + + def setUp(self) -> None: + self.clock = _FakeClock() + self.scheduler = Scheduler(clock=self.clock) + + def _names(self, claimed) -> list: + return sorted(collection.name for collection in claimed) + + def test_nothing_scheduled_is_quiet(self): + """Report no work and no delay before anything is scheduled.""" + self.assertIsNone(self.scheduler.next_delay()) + self.assertEqual(self.scheduler.claim_due(), ((), 0)) + + def test_everything_is_due_immediately_after_replace(self): + """Start collecting at once rather than after a first full interval.""" + self.scheduler.replace([_collection("adc"), _collection("pressure")]) + self.assertEqual(self.scheduler.next_delay(), 0.0) + claimed, skipped = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["adc", "pressure"]) + self.assertEqual(skipped, 0) + + def test_nothing_is_due_again_until_the_interval_passes(self): + """Hold a collection for its interval once it has run.""" + self.scheduler.replace([_collection("adc", interval_s=10.0)]) + self.scheduler.claim_due() + self.scheduler.release("adc") + + self.clock.advance(9.0) + self.assertEqual(self.scheduler.claim_due(), ((), 0)) + self.clock.advance(1.5) + claimed, _ = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["adc"]) + + def test_shorter_interval_comes_due_first(self): + """Order by due time, not by insertion.""" + self.scheduler.replace([_collection("slow", interval_s=60.0), + _collection("fast", interval_s=1.0)]) + self.scheduler.claim_due() + self.scheduler.release("slow") + self.scheduler.release("fast") + + self.clock.advance(2.0) + claimed, _ = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["fast"]) + + def test_tick_is_skipped_while_its_predecessor_runs(self): + """Count a skip rather than letting a slow peer overlap itself.""" + self.scheduler.replace([_collection("adc", interval_s=1.0)]) + self.scheduler.claim_due() # claimed, never released + + self.clock.advance(1.5) + claimed, skipped = self.scheduler.claim_due() + self.assertEqual(claimed, ()) + self.assertEqual(skipped, 1) + self.assertEqual(self.scheduler.in_flight, 1) + + def test_release_lets_the_next_tick_run(self): + """Resume a collection once its tick finishes.""" + self.scheduler.replace([_collection("adc", interval_s=1.0)]) + self.scheduler.claim_due() + self.scheduler.release("adc") + self.assertEqual(self.scheduler.in_flight, 0) + + self.clock.advance(1.5) + claimed, skipped = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["adc"]) + self.assertEqual(skipped, 0) + + def test_lag_records_how_late_a_tick_started(self): + """Report lateness, so a fleet that cannot keep up is visible.""" + collection = _collection("adc", interval_s=1.0) + self.scheduler.replace([collection]) + self.clock.advance(4.0) + self.scheduler.claim_due() + self.assertEqual(collection.lag_s, 4.0) + + def test_a_long_stall_does_not_queue_catch_up_ticks(self): + """Reschedule from now, so a stall cannot leave a burst that can only skip.""" + self.scheduler.replace([_collection("adc", interval_s=1.0)]) + self.scheduler.claim_due() + self.scheduler.release("adc") + + self.clock.advance(60.0) # a minute of stall on a 1s cadence + claimed, skipped = self.scheduler.claim_due() + self.assertEqual(len(claimed), 1) # one tick, not sixty + self.assertEqual(skipped, 0) + self.scheduler.release("adc") + self.assertEqual(self.scheduler.next_delay(), 1.0) + + def test_update_reschedules_only_its_own_collection(self): + """Change one cadence without re-seeding every other due time.""" + adc = _collection("adc", interval_s=10.0) + pressure = _collection("pressure", interval_s=10.0) + self.scheduler.replace([adc, pressure]) + self.scheduler.claim_due() + self.scheduler.release("adc") + self.scheduler.release("pressure") + + adc.config = dataclasses.replace(adc.config, interval_s=1.0) + self.scheduler.update(adc) + + claimed, _ = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["adc"]) # pressure still waits + + def test_replace_stops_scheduling_a_removed_collection(self): + """Drop a collection a reload removed.""" + self.scheduler.replace([_collection("adc"), _collection("pressure")]) + self.scheduler.claim_due() + self.scheduler.release("adc") + self.scheduler.release("pressure") + + self.scheduler.replace([_collection("adc")]) + self.assertEqual(self.scheduler.scheduled, ("adc",)) + claimed, _ = self.scheduler.claim_due() + self.assertEqual(self._names(claimed), ["adc"]) + + def test_replace_forgets_in_flight_names_it_dropped(self): + """Keep the in-flight set from leaking a collection that is gone.""" + self.scheduler.replace([_collection("adc"), _collection("pressure")]) + self.scheduler.claim_due() + self.assertEqual(self.scheduler.in_flight, 2) + + self.scheduler.replace([_collection("adc")]) + self.assertEqual(self.scheduler.in_flight, 1) + + def test_release_of_an_unknown_name_is_harmless(self): + """Tolerate a tick finishing for a collection a reload removed.""" + self.scheduler.release("nosuch") + self.assertEqual(self.scheduler.in_flight, 0) + + +if __name__ == "__main__": + unittest.main()