From 419ad76f3def9904610588e6062ffb27719d9338 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Thu, 17 Sep 2026 18:26:50 -0700 Subject: [PATCH] Report a failing sink, and back off a collection that keeps failing --- docs/source/keygrabber.md | 12 +++ libby/keygrabber/collection.py | 30 ++++++++ libby/keygrabber/daemon.py | 70 +++++++++++++++--- libby/keygrabber/scheduler.py | 5 +- libby/keygrabber/sink.py | 22 +++++- tests/test_keygrabber_backoff.py | 121 +++++++++++++++++++++++++++++++ tests/test_keygrabber_daemon.py | 33 ++++++++- 7 files changed, 279 insertions(+), 14 deletions(-) create mode 100644 tests/test_keygrabber_backoff.py diff --git a/docs/source/keygrabber.md b/docs/source/keygrabber.md index da5968f..aec470c 100644 --- a/docs/source/keygrabber.md +++ b/docs/source/keygrabber.md @@ -98,6 +98,14 @@ 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. +A collection that keeps failing reads less often. After three consecutive +failures its interval doubles each time, up to 32x, until a read succeeds. +Without that, a daemon that is down would be retried on its configured cadence +indefinitely, logging an error every time and rewriting `lasterror` with it. +Logging follows the same shape: the first failures are reported, then the +backoff is announced once, then it stays quiet until the peer answers again. +The cap means a peer that comes back is picked up within a bounded time. + ## Control keywords The keygrabber is itself a peer, so its cadence and health are reachable with @@ -136,6 +144,10 @@ 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. +A failing sink is visible from three keywords together: `isconnected` goes +false, `queuedepth` climbs as batches wait, and `pointswritten` stops moving. +`writeerrors` counts every failed attempt including retries. + ### reload `reload` re-reads the config file and applies it to the running collections, diff --git a/libby/keygrabber/collection.py b/libby/keygrabber/collection.py index ea2e742..f4eda6b 100644 --- a/libby/keygrabber/collection.py +++ b/libby/keygrabber/collection.py @@ -13,6 +13,14 @@ BULK_READ_SERVICE = "keys.read" +# Failures tolerated at the configured cadence before a collection starts +# reading less often. A peer restarting should not trigger a backoff. +FAILURES_BEFORE_BACKOFF = 3 + +# Ceiling on how far the interval is stretched while a peer stays silent, so a +# recovered peer is picked up again within a bounded time +MAX_BACKOFF_MULTIPLIER = 32 + @dataclass(frozen=True) class TickResult: @@ -45,6 +53,7 @@ def __init__( self.enabled = True self.last_sample: Optional[datetime] = None self.lag_s = 0.0 + self.consecutive_failures = 0 self._clock = clock self._names: Tuple[str, ...] = () self._bulk_read = False @@ -75,6 +84,27 @@ def invalidate(self) -> None: """Force the next tick to resolve again, after a config change.""" self._resolved_at = None + def note_failure(self) -> None: + """Record a tick that read nothing.""" + self.consecutive_failures += 1 + + def note_success(self) -> None: + """Record a tick that read something, ending any backoff.""" + self.consecutive_failures = 0 + + def backoff_interval_s(self) -> float: + """Return the interval to use next, stretched while the peer fails. + + A peer that is down would otherwise be retried, and logged about, on + its configured cadence indefinitely. Backing off keeps a dead daemon + from dominating both the logs and the read budget, while the cap keeps + a recovered one from waiting long to be noticed. + """ + if self.consecutive_failures < FAILURES_BEFORE_BACKOFF: + return self.config.interval_s + overshoot = self.consecutive_failures - FAILURES_BEFORE_BACKOFF + 1 + return self.config.interval_s * min(2 ** overshoot, MAX_BACKOFF_MULTIPLIER) + 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 a65c1c1..c66224f 100644 --- a/libby/keygrabber/daemon.py +++ b/libby/keygrabber/daemon.py @@ -14,7 +14,7 @@ from ..daemon import LibbyDaemon from ..errors import LibbyError from ..libby import Libby -from .collection import Collection +from .collection import FAILURES_BEFORE_BACKOFF, Collection from .config import TIMEOUT_HEADROOM, KeygrabberConfig, build_sink, parse_config from .scheduler import Scheduler from .sink import RetryingWriter, Sample, Sink @@ -73,6 +73,12 @@ def add(self, **deltas: int) -> None: for name, delta in deltas.items(): setattr(self, name, getattr(self, name) + delta) + def set(self, **values: int) -> None: + """Set one or more counters to an absolute value.""" + with self._lock: + for name, value in values.items(): + setattr(self, name, value) + # Coordinates config, client, sink, collections and three kinds of thread; # the attribute count is collaborators, not state @@ -364,13 +370,38 @@ 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) + collection.last_sample = datetime.now(timezone.utc) + self._note_success(collection) + elif collection.keyword_count: + self._note_failure(collection, "every read failed") except LibbyError as exc: self.counters.add(read_errors=1) - self.logger.error("collection %s read failed: %s", collection.name, exc) + self._note_failure(collection, str(exc)) finally: self._scheduler.release(collection.name) + def _note_failure(self, collection: Collection, reason: str) -> None: + """Count a failed tick, and log only while that is still news. + + A peer that stays down would otherwise produce an error every interval + for as long as it is down, which buries everything else and keeps + rewriting ``lasterror``. + """ + collection.note_failure() + failures = collection.consecutive_failures + if failures < FAILURES_BEFORE_BACKOFF: + self.logger.error("collection %s failed: %s", collection.name, reason) + elif failures == FAILURES_BEFORE_BACKOFF: + self.logger.error( + "collection %s has failed %d times, backing off to %.0fs: %s", + collection.name, failures, collection.backoff_interval_s(), reason) + + def _note_success(self, collection: Collection) -> None: + """Clear a collection's backoff, saying so if it had been failing.""" + if collection.consecutive_failures >= FAILURES_BEFORE_BACKOFF: + self.logger.info("collection %s is answering again", collection.name) + collection.note_success() + def _enqueue(self, samples: Sequence[Sample]) -> None: try: self._queue.put_nowait(tuple(samples)) @@ -417,11 +448,12 @@ def _write(self, batch: Tuple[Sample, ...]) -> 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) + # A sink raising something other than SinkWriteError is a bug in + # that sink; the retry writer turns the expected failure into a + # queued batch instead + self.logger.error("sink write raised: %s", exc) + self._track_sink_health(writer) def _flush(self) -> None: writer = self._writer @@ -430,9 +462,27 @@ def _flush(self) -> None: try: 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) + self.logger.error("sink retry raised: %s", exc) + self._track_sink_health(writer) + + def _track_sink_health(self, writer: RetryingWriter) -> None: + """Mirror the writer's view, and log only when it changes. + + ``RetryingWriter.write`` queues a failed batch and returns 0 rather + than raising, which is what makes retrying possible. It also means a + failure cannot be noticed from an exception here, so health and the + error count are read back from the writer instead. + """ + self.counters.set(write_errors=writer.failed_attempts) + healthy = writer.healthy + if healthy == self._sink_healthy: + return + self._sink_healthy = healthy + if healthy: + self.logger.info("sink writes are succeeding again") + else: + self.logger.error("sink writes are failing, %d batch(es) queued", + writer.queue_depth) def _drain(self) -> None: """Write whatever is still queued, under a deadline.""" diff --git a/libby/keygrabber/scheduler.py b/libby/keygrabber/scheduler.py index f0ec5d4..bfe7342 100644 --- a/libby/keygrabber/scheduler.py +++ b/libby/keygrabber/scheduler.py @@ -77,9 +77,10 @@ def claim_due(self) -> Tuple[Tuple[Collection, ...], int]: 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. + # skip. Cadence drifts by the loop's own latency instead, and + # stretches while the peer is failing. heapq.heappush( - self._due, (now + collection.config.interval_s, name)) + self._due, (now + collection.backoff_interval_s(), name)) return tuple(claimed), skipped def update(self, collection: Collection) -> None: diff --git a/libby/keygrabber/sink.py b/libby/keygrabber/sink.py index 19b98f9..c660021 100644 --- a/libby/keygrabber/sink.py +++ b/libby/keygrabber/sink.py @@ -78,7 +78,9 @@ def close(self) -> None: """Release the connection.""" -class RetryingWriter: +# The extra attributes are the queue, the backoff and the tallies it reports, +# each one value rather than hidden state +class RetryingWriter: # pylint: disable=too-many-instance-attributes """Wraps a sink, holding failed batches in a bounded queue for retry. Retrying is backend-independent, so it lives here rather than inside any @@ -102,6 +104,7 @@ def __init__( self._clock = clock self._pending: Deque[Tuple[Sample, ...]] = deque() self._failures = 0 + self._failed_attempts = 0 self._retry_at = 0.0 self._dropped_batches = 0 @@ -115,6 +118,22 @@ def dropped_batches(self) -> int: """Number of batches discarded because the queue was full.""" return self._dropped_batches + @property + def healthy(self) -> bool: + """Whether the most recent attempt on the sink succeeded. + + Worth asking for, because :meth:`write` deliberately does not raise + when a batch fails: it queues it and returns 0, which is what makes + retrying possible and what stops a caller learning about the failure + from an exception. + """ + return self._failures == 0 + + @property + def failed_attempts(self) -> int: + """Total attempts on the sink that failed, including retries.""" + return self._failed_attempts + def write(self, samples: Sequence[Sample]) -> int: """Write a batch now, or queue it if the sink is in backoff.""" batch = tuple(samples) @@ -176,6 +195,7 @@ def _backoff_elapsed(self) -> bool: def _arm_backoff(self) -> None: self._failures += 1 + self._failed_attempts += 1 delay = min(self._policy.base_backoff_s * (2 ** (self._failures - 1)), self._policy.max_backoff_s) self._retry_at = self._clock() + delay diff --git a/tests/test_keygrabber_backoff.py b/tests/test_keygrabber_backoff.py new file mode 100644 index 0000000..3c6da0e --- /dev/null +++ b/tests/test_keygrabber_backoff.py @@ -0,0 +1,121 @@ +"""Unit tests for a collection backing off while its peer stays silent.""" +from __future__ import annotations + +import unittest + +from libby.keygrabber import Collection +from libby.keygrabber.collection import ( + FAILURES_BEFORE_BACKOFF, + MAX_BACKOFF_MULTIPLIER, +) +from libby.keygrabber.config import CollectionConfig +from libby.keygrabber.scheduler import Scheduler + +INTERVAL_S = 10.0 + + +def _collection() -> Collection: + return Collection(CollectionConfig( + name="adc", group="hsfei", daemon="adc", 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 BackoffIntervalTests(unittest.TestCase): + """How a collection stretches its own cadence while failing.""" + + def setUp(self) -> None: + self.collection = _collection() + + def test_a_healthy_collection_uses_its_configured_cadence(self): + """Leave the interval alone when nothing has failed.""" + self.assertEqual(self.collection.backoff_interval_s(), INTERVAL_S) + + def test_early_failures_do_not_slow_it_down(self): + """Tolerate a peer restarting without changing its cadence.""" + for _ in range(FAILURES_BEFORE_BACKOFF - 1): + self.collection.note_failure() + self.assertEqual(self.collection.backoff_interval_s(), INTERVAL_S) + + def test_the_interval_grows_once_failures_persist(self): + """Read a persistently dead peer less often.""" + for _ in range(FAILURES_BEFORE_BACKOFF): + self.collection.note_failure() + first = self.collection.backoff_interval_s() + self.assertGreater(first, INTERVAL_S) + + self.collection.note_failure() + self.assertGreater(self.collection.backoff_interval_s(), first) + + def test_the_interval_is_capped(self): + """Keep a recovered peer from waiting indefinitely to be noticed.""" + for _ in range(100): + self.collection.note_failure() + self.assertEqual(self.collection.backoff_interval_s(), + INTERVAL_S * MAX_BACKOFF_MULTIPLIER) + + def test_one_success_clears_the_backoff(self): + """Return to the configured cadence as soon as a read works.""" + for _ in range(20): + self.collection.note_failure() + self.collection.note_success() + self.assertEqual(self.collection.consecutive_failures, 0) + self.assertEqual(self.collection.backoff_interval_s(), INTERVAL_S) + + +class SchedulerHonoursBackoffTests(unittest.TestCase): + """The scheduler reads the stretched interval, not the configured one.""" + + def test_a_failing_collection_is_scheduled_further_out(self): + """Space out the next tick of a collection that keeps failing.""" + clock = _FakeClock() + scheduler = Scheduler(clock=clock) + collection = _collection() + scheduler.replace([collection]) + + scheduler.claim_due() + scheduler.release(collection.name) + self.assertEqual(scheduler.next_delay(), INTERVAL_S) + + for _ in range(FAILURES_BEFORE_BACKOFF): + collection.note_failure() + clock.advance(INTERVAL_S) + scheduler.claim_due() + scheduler.release(collection.name) + self.assertGreater(scheduler.next_delay(), INTERVAL_S) + + def test_recovery_restores_the_configured_cadence(self): + """Go back to the normal interval on the first successful read.""" + clock = _FakeClock() + scheduler = Scheduler(clock=clock) + collection = _collection() + scheduler.replace([collection]) + + for _ in range(FAILURES_BEFORE_BACKOFF + 2): + collection.note_failure() + scheduler.claim_due() + scheduler.release(collection.name) + self.assertGreater(scheduler.next_delay(), INTERVAL_S) + + collection.note_success() + clock.advance(scheduler.next_delay()) + scheduler.claim_due() + scheduler.release(collection.name) + self.assertEqual(scheduler.next_delay(), INTERVAL_S) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_keygrabber_daemon.py b/tests/test_keygrabber_daemon.py index db33819..2f88ea3 100644 --- a/tests/test_keygrabber_daemon.py +++ b/tests/test_keygrabber_daemon.py @@ -17,7 +17,7 @@ from libby import Client, KeywordError from libby.daemon import LibbyDaemon -from libby.keygrabber import KeygrabberDaemon, Sample +from libby.keygrabber import KeygrabberDaemon, Sample, SinkWriteError from libby.rabbitmq_transport import RabbitMQTransport SETTLE_TIMEOUT_S = 15.0 @@ -102,6 +102,14 @@ def snapshot(self) -> List[Sample]: return list(self.samples) +class _FailingSink(_RecordingSink): + """Sink that refuses every write, as an unreachable database does.""" + + def write(self, samples: Sequence[Sample]) -> int: + """Reject the batch the way a dead backend does.""" + raise SinkWriteError("backend unreachable") + + class _TestKeygrabber(KeygrabberDaemon): """Keygrabber writing to a recording sink instead of InfluxDB. @@ -289,6 +297,29 @@ def test_reload_without_a_config_file_is_refused(self): with self.assertRaises(KeywordError): client.set(self._keyword("reload"), 1) + def test_a_dead_sink_is_reported_not_hidden(self): + """Report an unwritable sink, rather than looking healthy. + + ``RetryingWriter.write`` queues a failed batch and returns 0 + instead of raising, so health has to be read back from the writer. + Getting this wrong showed a live daemon with queued batches, + nothing written, and ``isconnected`` still true. + """ + self.grabber.sink = _FailingSink() + self.grabber.start() + deadline = time.monotonic() + SETTLE_TIMEOUT_S + while time.monotonic() < deadline: + if self.grabber.counters.write_errors: + break + time.sleep(0.05) + + with self._grabber_client() as client: + self.assertFalse(client.get(self._keyword("isconnected"))) + self.assertEqual(client.get(self._keyword("pointswritten")), 0) + self.assertGreater(client.get(self._keyword("writeerrors")), 0) + self.assertGreater(client.get(self._keyword("queuedepth")), 0) + self.assertIsNotNone(client.get(self._keyword("lasterror"))) + def test_repeats_on_the_configured_cadence(self): """Read again on the next interval rather than once at startup.""" self.grabber.start()