Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions docs/source/keygrabber.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down
30 changes: 30 additions & 0 deletions libby/keygrabber/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down
70 changes: 60 additions & 10 deletions libby/keygrabber/daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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
Expand All @@ -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."""
Expand Down
5 changes: 3 additions & 2 deletions libby/keygrabber/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
22 changes: 21 additions & 1 deletion libby/keygrabber/sink.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand Down
121 changes: 121 additions & 0 deletions tests/test_keygrabber_backoff.py
Original file line number Diff line number Diff line change
@@ -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()
Loading
Loading