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
11 changes: 11 additions & 0 deletions docs/source/api/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,17 @@ Generated from docstrings in the `libby` package.
:show-inheritance:
```

## Keygrabber

The InfluxDB sink is omitted here: it imports its client at module scope, and
the docs build does not install the optional `influxdb` extra.

```{eval-rst}
.. automodule:: libby.keygrabber.sink
:members:
:show-inheritance:
```

## Responses and errors

```{eval-rst}
Expand Down
24 changes: 24 additions & 0 deletions libby/keygrabber/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
"""Keygrabber: poll keywords from libby peers and store them as time series.

Importing this package does not pull in any database client. ``InfluxSink``
lives in :mod:`libby.keygrabber.influx` and needs the optional dependency
(``pip install libby[influxdb]``), so a deployment that only uses another
backend never has to install it.
"""
from .sink import (
RetryingWriter,
RetryPolicy,
Sample,
Sink,
SinkError,
SinkWriteError,
)

__all__ = [
"RetryingWriter",
"RetryPolicy",
"Sample",
"Sink",
"SinkError",
"SinkWriteError",
]
123 changes: 123 additions & 0 deletions libby/keygrabber/influx.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
"""InfluxDB 2.x sink.

Schema is one measurement per keyword name, tagged by ``group``, ``peer`` and
``units``, with a single ``value`` field. A keyword name carries one type across
peers, so field types stay consistent, and a Grafana query is a measurement
plus a ``peer`` tag filter.

Needs the optional dependency: ``pip install libby[influxdb]``.
"""
from __future__ import annotations

from dataclasses import dataclass, field
from typing import Optional, Sequence

from influxdb_client import InfluxDBClient, Point, WritePrecision
from influxdb_client.client.write_api import SYNCHRONOUS

from .sink import Sample, SinkWriteError, Value

# Influx drops an empty tag value, which would split one keyword into two
# series depending on whether it declared units. A literal keeps every point
# on the same series.
NO_UNITS = "none"

DEFAULT_TIMEOUT_MS = 10_000


def field_value(value: Value) -> Value:
"""Coerce a keyword value to the field type Influx should store.

Ints become floats so that one peer reporting ``0`` and another ``0.5``
for the same keyword cannot collide as int against float and have the
write rejected. ``bool`` is checked first, being a subclass of ``int``.
"""
if isinstance(value, bool):
return value
if isinstance(value, (int, float)):
return float(value)
return str(value)


def to_point(sample: Sample) -> Optional[Point]:
"""Map a sample to a point, or ``None`` if Influx cannot store it.

A null value is skipped rather than raised: ``lasterror`` is nullable on
every daemon and is ``None`` most of the time, so a collection reading a
whole peer would otherwise fail on every tick.
"""
if sample.value is None:
return None
return (
Point(sample.keyword)
.tag("group", sample.group)
.tag("peer", sample.peer)
.tag("units", sample.units or NO_UNITS)
.field("value", field_value(sample.value))
.time(sample.timestamp, WritePrecision.NS)
)


@dataclass(frozen=True)
class InfluxConfig:
"""Connection settings for one InfluxDB 2.x bucket.

``token`` is kept out of ``repr`` so a traceback or a logged config cannot
leak it.
"""

url: str
org: str
bucket: str
token: str = field(repr=False)
timeout_ms: int = DEFAULT_TIMEOUT_MS


class InfluxSink:
"""Writes samples to an InfluxDB 2.x bucket."""

def __init__(self, config: InfluxConfig) -> None:
self._config = config
self._client: Optional[InfluxDBClient] = None
self._write_api = None

def connect(self) -> None:
"""Open the client, replacing any existing one."""
self.close()
self._client = InfluxDBClient(url=self._config.url,
token=self._config.token,
org=self._config.org,
timeout=self._config.timeout_ms)
self._write_api = self._client.write_api(write_options=SYNCHRONOUS)

def is_connected(self) -> bool:
"""Return whether the server answers a ping."""
if self._client is None:
return False
try:
return bool(self._client.ping())
except Exception: # pylint: disable=broad-exception-caught
return False

def write(self, samples: Sequence[Sample]) -> int:
"""Write a batch as one request, skipping samples Influx cannot store."""
if self._write_api is None:
raise SinkWriteError("influx sink is not connected")
points = [point for point in map(to_point, samples) if point is not None]
if not points:
return 0
try:
self._write_api.write(bucket=self._config.bucket,
org=self._config.org, record=points)
# The client surfaces API, HTTP and socket errors with no common base,
# and the caller's contract is a single retryable error
except Exception as exc: # pylint: disable=broad-exception-caught
raise SinkWriteError(f"influx write failed: {exc}") from exc
return len(points)

def close(self) -> None:
"""Release the client; safe to call when never connected."""
client, self._client = self._client, None
self._write_api = None
if client is not None:
client.close()
185 changes: 185 additions & 0 deletions libby/keygrabber/sink.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
"""The sample the keygrabber collects, and the sink contract it writes to.

A ``Sample`` carries nothing backend-specific: turning one into measurements,
fields, tags or columns is the sink's job, so a second backend is a new sink
rather than a change to the collector.
"""
from __future__ import annotations

import time
from collections import deque
from dataclasses import dataclass
from datetime import datetime
from typing import Callable, Deque, Optional, Protocol, Sequence, Tuple, Union

from ..errors import LibbyError

Value = Union[bool, int, float, str, None]

DEFAULT_MAX_BATCHES = 64
DEFAULT_BASE_BACKOFF_S = 1.0
DEFAULT_MAX_BACKOFF_S = 60.0


class SinkError(LibbyError):
"""A sink could not be reached, configured, or written to."""


class SinkWriteError(SinkError):
"""A batch of samples could not be written and may be worth retrying."""


@dataclass(frozen=True)
class Sample:
"""One keyword's value, read at one instant."""

keyword: str
group: str
peer: str
value: Value
units: Optional[str]
timestamp: datetime


@dataclass(frozen=True)
class RetryPolicy:
"""Bounds on how a :class:`RetryingWriter` queues and re-attempts batches."""

max_batches: int = DEFAULT_MAX_BATCHES
base_backoff_s: float = DEFAULT_BASE_BACKOFF_S
max_backoff_s: float = DEFAULT_MAX_BACKOFF_S

def __post_init__(self) -> None:
if self.max_batches < 1:
raise ValueError("max_batches must be at least 1")
if self.base_backoff_s <= 0:
raise ValueError("base_backoff_s must be positive")


class Sink(Protocol):
"""Destination for collected samples.

Implementations own their own schema. ``write`` returns how many samples it
actually stored, which can be fewer than it was given when a backend cannot
represent some of them, and raises :class:`SinkWriteError` when the batch
failed and should be retried.
"""

def connect(self) -> None:
"""Open the connection, or reopen it after a failure."""

def is_connected(self) -> bool:
"""Return whether the backend is currently reachable."""

def write(self, samples: Sequence[Sample]) -> int:
"""Store a batch and return how many samples were written."""

def close(self) -> None:
"""Release the connection."""


class RetryingWriter:
"""Wraps a sink, holding failed batches in a bounded queue for retry.

Retrying is backend-independent, so it lives here rather than inside any
one sink. The queue is bounded and drops its oldest batch when full, so a
database that stays down cannot grow the daemon's memory without limit.

Nothing here starts a thread and nothing sleeps: the caller decides when to
call :meth:`flush_due`, and ``clock`` is injectable so backoff is testable
without waiting for it.
"""

def __init__(
self,
sink: Sink,
*,
policy: Optional[RetryPolicy] = None,
clock: Callable[[], float] = time.monotonic,
) -> None:
self._sink = sink
self._policy = policy or RetryPolicy()
self._clock = clock
self._pending: Deque[Tuple[Sample, ...]] = deque()
self._failures = 0
self._retry_at = 0.0
self._dropped_batches = 0

@property
def queue_depth(self) -> int:
"""Number of batches waiting to be retried."""
return len(self._pending)

@property
def dropped_batches(self) -> int:
"""Number of batches discarded because the queue was full."""
return self._dropped_batches

def write(self, samples: Sequence[Sample]) -> int:
"""Write a batch now, or queue it if the sink is in backoff."""
batch = tuple(samples)
if not batch:
return 0
# Queue behind an existing backlog rather than overtaking it, and do
# not probe a sink that is still inside its backoff window
if self._pending or not self._backoff_elapsed():
self._enqueue(batch)
return 0
return self._attempt(batch)

def flush_due(self) -> int:
"""Retry queued batches, oldest first, once backoff has elapsed."""
if not self._pending or not self._backoff_elapsed():
return 0
written = 0
while self._pending:
try:
written += self._sink.write(self._pending[0])
except SinkWriteError:
self._arm_backoff()
return written
self._pending.popleft()
self._reset_backoff()
return written

def connect(self) -> None:
"""Reconnect the wrapped sink and clear its backoff."""
self._sink.connect()
self._reset_backoff()

def is_connected(self) -> bool:
"""Return whether the wrapped sink is reachable."""
return self._sink.is_connected()

def close(self) -> None:
"""Release the wrapped sink."""
self._sink.close()

def _attempt(self, batch: Tuple[Sample, ...]) -> int:
try:
written = self._sink.write(batch)
except SinkWriteError:
self._enqueue(batch)
self._arm_backoff()
return 0
self._reset_backoff()
return written

def _enqueue(self, batch: Tuple[Sample, ...]) -> None:
if len(self._pending) >= self._policy.max_batches:
self._pending.popleft()
self._dropped_batches += 1
self._pending.append(batch)

def _backoff_elapsed(self) -> bool:
return self._clock() >= self._retry_at

def _arm_backoff(self) -> None:
self._failures += 1
delay = min(self._policy.base_backoff_s * (2 ** (self._failures - 1)),
self._policy.max_backoff_s)
self._retry_at = self._clock() + delay

def _reset_backoff(self) -> None:
self._failures = 0
self._retry_at = 0.0
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ dependencies = [
]

[project.optional-dependencies]
influxdb = [
"influxdb-client>=1.40",
]
docs = [
"sphinx>=7.2",
"shibuya>=2024.1.21",
Expand Down
Loading
Loading