From 1037a79a126f4ba20fb923b7a94b9609086bbfdf Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:40:07 -0700 Subject: [PATCH] Add the keygrabber sample, sink contract, retry queue and InfluxDB sink --- docs/source/api/index.md | 11 ++ libby/keygrabber/__init__.py | 24 +++++ libby/keygrabber/influx.py | 123 +++++++++++++++++++++++ libby/keygrabber/sink.py | 185 ++++++++++++++++++++++++++++++++++ pyproject.toml | 3 + tests/test_influx_sink.py | 167 +++++++++++++++++++++++++++++++ tests/test_sink.py | 189 +++++++++++++++++++++++++++++++++++ tox.ini | 3 +- 8 files changed, 704 insertions(+), 1 deletion(-) create mode 100644 libby/keygrabber/__init__.py create mode 100644 libby/keygrabber/influx.py create mode 100644 libby/keygrabber/sink.py create mode 100644 tests/test_influx_sink.py create mode 100644 tests/test_sink.py diff --git a/docs/source/api/index.md b/docs/source/api/index.md index 998af93..6746733 100644 --- a/docs/source/api/index.md +++ b/docs/source/api/index.md @@ -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} diff --git a/libby/keygrabber/__init__.py b/libby/keygrabber/__init__.py new file mode 100644 index 0000000..064b9c0 --- /dev/null +++ b/libby/keygrabber/__init__.py @@ -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", +] diff --git a/libby/keygrabber/influx.py b/libby/keygrabber/influx.py new file mode 100644 index 0000000..ff366b1 --- /dev/null +++ b/libby/keygrabber/influx.py @@ -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() diff --git a/libby/keygrabber/sink.py b/libby/keygrabber/sink.py new file mode 100644 index 0000000..19b98f9 --- /dev/null +++ b/libby/keygrabber/sink.py @@ -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 diff --git a/pyproject.toml b/pyproject.toml index e0c077e..616e507 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -22,6 +22,9 @@ dependencies = [ ] [project.optional-dependencies] +influxdb = [ + "influxdb-client>=1.40", +] docs = [ "sphinx>=7.2", "shibuya>=2024.1.21", diff --git a/tests/test_influx_sink.py b/tests/test_influx_sink.py new file mode 100644 index 0000000..0fb7401 --- /dev/null +++ b/tests/test_influx_sink.py @@ -0,0 +1,167 @@ +"""Unit tests for the InfluxDB sink's sample mapping and error contract. + +Needs the optional dependency (``pip install libby[influxdb]``); tox installs +it, so these run in CI rather than skipping. +""" +from __future__ import annotations + +import importlib.util +import unittest +from datetime import datetime, timezone +from typing import Any, List, Optional + +from libby.keygrabber import Sample, SinkWriteError + +HAS_INFLUX = importlib.util.find_spec("influxdb_client") is not None + +# Every class below is skipUnless-guarded on HAS_INFLUX, which pylint +# cannot see through +# pylint: disable=possibly-used-before-assignment + +if HAS_INFLUX: + from libby.keygrabber.influx import (NO_UNITS, InfluxConfig, InfluxSink, + field_value, to_point) + + +def _sample( + keyword: str = "positionvalue", + value: Any = 7.5, + units: Optional[str] = "mm", +) -> Sample: + return Sample(keyword=keyword, group="hsfei", peer="adc", value=value, + units=units, + timestamp=datetime(2026, 1, 1, 12, 0, tzinfo=timezone.utc)) + + +class _FakeWriteApi: # pylint: disable=too-few-public-methods + """Records write calls, or raises to exercise the error contract.""" + + def __init__(self, error: Optional[Exception] = None) -> None: + self.records: List[Any] = [] + self._error = error + + def write(self, bucket: str, org: str, record: Any) -> None: + """Record the batch, or raise the configured error.""" + if self._error is not None: + raise self._error + self.records.append((bucket, org, record)) + + +@unittest.skipUnless(HAS_INFLUX, "influxdb-client not installed") +class FieldValueTests(unittest.TestCase): + """Coercion of keyword values to Influx field types.""" + + def test_int_becomes_float(self): + """Avoid an int/float type conflict between peers for one keyword.""" + coerced = field_value(12) + self.assertIsInstance(coerced, float) + self.assertEqual(coerced, 12.0) + + def test_bool_stays_bool(self): + """Keep bools as bools, despite bool being a subclass of int.""" + self.assertIs(field_value(True), True) + self.assertIsInstance(field_value(False), bool) + + def test_float_stays_float(self): + """Pass a float through unchanged.""" + self.assertEqual(field_value(7.5), 7.5) + + def test_string_stays_string(self): + """Store a string keyword as a string field.""" + self.assertEqual(field_value("Ready"), "Ready") + + +@unittest.skipUnless(HAS_INFLUX, "influxdb-client not installed") +class ToPointTests(unittest.TestCase): + """Mapping a sample onto the measurement-per-keyword schema.""" + + def test_measurement_is_the_keyword_name(self): + """Use the keyword name as the measurement, with group/peer as tags.""" + line = to_point(_sample()).to_line_protocol() + self.assertTrue(line.startswith("positionvalue,")) + self.assertIn("group=hsfei", line) + self.assertIn("peer=adc", line) + self.assertIn("units=mm", line) + self.assertIn("value=7.5", line) + + def test_null_value_is_skipped(self): + """Skip a null rather than raising: lasterror is null most of the time.""" + self.assertIsNone(to_point(_sample(keyword="lasterror", value=None))) + + def test_missing_units_get_a_placeholder(self): + """Keep a unitless keyword on one series instead of dropping the tag.""" + line = to_point(_sample(units=None)).to_line_protocol() + self.assertIn(f"units={NO_UNITS}", line) + + def test_timestamp_is_carried_at_nanosecond_precision(self): + """Stamp the point with the sample's own read time.""" + line = to_point(_sample()).to_line_protocol() + expected_ns = int(_sample().timestamp.timestamp() * 1_000_000_000) + self.assertTrue(line.endswith(str(expected_ns))) + + def test_string_value_is_quoted_by_the_client(self): + """Let the client escape a string field rather than hand-rolling it.""" + line = to_point(_sample(keyword="lasterror", + value='failed: "x", y', units=None)).to_line_protocol() + self.assertIn('value="failed: \\"x\\", y"', line) + + +@unittest.skipUnless(HAS_INFLUX, "influxdb-client not installed") +class InfluxSinkWriteTests(unittest.TestCase): + """The sink's write contract, without a server.""" + + # The write api is injected in place of connect() so these need no server + # pylint: disable=protected-access + + def setUp(self) -> None: + self.sink = InfluxSink(InfluxConfig( + url="http://localhost:8086", org="hispec", + bucket="telemetry", token="unused")) + + def test_writes_one_request_for_the_whole_batch(self): + """Send a batch as a single request and report the count written.""" + api = _FakeWriteApi() + self.sink._write_api = api + written = self.sink.write([_sample(), _sample(keyword="isconnected", + value=True, units=None)]) + self.assertEqual(written, 2) + self.assertEqual(len(api.records), 1) + self.assertEqual(len(api.records[0][2]), 2) + + def test_nulls_are_not_counted_as_written(self): + """Report only the samples that reached the backend.""" + api = _FakeWriteApi() + self.sink._write_api = api + written = self.sink.write([_sample(), + _sample(keyword="lasterror", value=None)]) + self.assertEqual(written, 1) + + def test_all_null_batch_makes_no_request(self): + """Skip the request entirely when nothing is storable.""" + api = _FakeWriteApi() + self.sink._write_api = api + self.assertEqual(self.sink.write([_sample(value=None)]), 0) + self.assertEqual(api.records, []) + + def test_backend_failure_becomes_a_sink_write_error(self): + """Translate the client's own exceptions into one retryable error.""" + self.sink._write_api = _FakeWriteApi(error=OSError("connection refused")) + with self.assertRaises(SinkWriteError): + self.sink.write([_sample()]) + + def test_writing_before_connect_is_a_sink_write_error(self): + """Fail a write on an unconnected sink the same retryable way.""" + with self.assertRaises(SinkWriteError): + self.sink.write([_sample()]) + + def test_is_connected_is_false_before_connect(self): + """Report not connected rather than raising.""" + self.assertFalse(self.sink.is_connected()) + + def test_close_without_connect_is_safe(self): + """Allow teardown of a sink that never opened.""" + self.sink.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_sink.py b/tests/test_sink.py new file mode 100644 index 0000000..b5e6509 --- /dev/null +++ b/tests/test_sink.py @@ -0,0 +1,189 @@ +"""Unit tests for RetryingWriter: queueing, backoff, bounds and ordering. + +The clock is injected, so nothing here waits for a real backoff window. +""" +from __future__ import annotations + +import unittest +from datetime import datetime, timezone +from typing import List, Sequence, Tuple + +from libby.keygrabber import RetryingWriter, RetryPolicy, Sample, SinkWriteError + + +def _sample(keyword: str = "positionvalue", value: float = 1.0) -> Sample: + return Sample(keyword=keyword, group="hsfei", peer="adc", value=value, + units="mm", timestamp=datetime(2026, 1, 1, tzinfo=timezone.utc)) + + +class _FakeClock: + """Manually advanced clock, so backoff is exercised without sleeping.""" + + 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 _FakeSink: + """Sink that records what it was given and fails on demand.""" + + def __init__(self) -> None: + self.batches: List[Tuple[Sample, ...]] = [] + self.attempts = 0 + self.failing = False + self.connects = 0 + self.closed = False + + def connect(self) -> None: + """Count reconnections.""" + self.connects += 1 + + def is_connected(self) -> bool: + """Report the inverse of the failing flag.""" + return not self.failing + + def write(self, samples: Sequence[Sample]) -> int: + """Record the batch, or raise while failing.""" + self.attempts += 1 + if self.failing: + raise SinkWriteError("backend down") + self.batches.append(tuple(samples)) + return len(samples) + + def close(self) -> None: + """Mark the sink closed.""" + self.closed = True + + +class RetryingWriterTests(unittest.TestCase): + """Behaviour of the retry queue in front of a sink.""" + + def setUp(self) -> None: + self.clock = _FakeClock() + self.sink = _FakeSink() + self.writer = RetryingWriter( + self.sink, + policy=RetryPolicy(max_batches=3, base_backoff_s=1.0, max_backoff_s=8.0), + clock=self.clock) + + def test_healthy_write_goes_straight_through(self): + """Pass a batch to the sink and report what it wrote.""" + self.assertEqual(self.writer.write([_sample(), _sample()]), 2) + self.assertEqual(len(self.sink.batches), 1) + self.assertEqual(self.writer.queue_depth, 0) + + def test_empty_batch_is_not_written(self): + """Skip an empty batch rather than making a pointless request.""" + self.assertEqual(self.writer.write([]), 0) + self.assertEqual(self.sink.batches, []) + + def test_failed_write_is_queued_not_lost(self): + """Hold a failed batch for retry and report nothing written.""" + self.sink.failing = True + self.assertEqual(self.writer.write([_sample()]), 0) + self.assertEqual(self.writer.queue_depth, 1) + + def test_queued_batch_is_written_once_backoff_elapses(self): + """Retry the backlog after the backoff window, not before.""" + self.sink.failing = True + self.writer.write([_sample()]) + self.sink.failing = False + + self.assertEqual(self.writer.flush_due(), 0) # still inside backoff + self.assertEqual(self.writer.queue_depth, 1) + + self.clock.advance(1.0) + self.assertEqual(self.writer.flush_due(), 1) + self.assertEqual(self.writer.queue_depth, 0) + + def test_backoff_grows_and_is_capped(self): + """Wait longer between retries as failures repeat, up to the ceiling.""" + self.sink.failing = True + self.writer.write([_sample()]) # first failure arms the backoff + + waits = [] + for _ in range(5): + attempts_before = self.sink.attempts + waited = 0.0 + # Measure when the writer next touches the sink at all, since a + # retry that fails is still a retry + while self.sink.attempts == attempts_before and waited < 30.0: + self.clock.advance(0.5) + waited += 0.5 + self.writer.flush_due() + waits.append(waited) + self.assertEqual(waits[:4], [1.0, 2.0, 4.0, 8.0]) + self.assertEqual(waits[4], 8.0) # capped at max_backoff_s + + def test_queue_drops_oldest_when_full(self): + """Bound the queue so a dead backend cannot grow memory without limit.""" + self.sink.failing = True + for index in range(5): + self.writer.write([_sample(value=float(index))]) + self.assertEqual(self.writer.queue_depth, 3) + self.assertEqual(self.writer.dropped_batches, 2) + + self.sink.failing = False + self.clock.advance(100.0) + self.writer.flush_due() + # The two oldest went, the three newest survived in order + written = [batch[0].value for batch in self.sink.batches] + self.assertEqual(written, [2.0, 3.0, 4.0]) + + def test_write_queues_behind_an_existing_backlog(self): + """Keep ordering by not overtaking a backlog with a fresh batch.""" + self.sink.failing = True + self.writer.write([_sample(value=1.0)]) + self.sink.failing = False + self.assertEqual(self.writer.write([_sample(value=2.0)]), 0) + self.assertEqual(self.writer.queue_depth, 2) + + self.clock.advance(1.0) + self.writer.flush_due() + self.assertEqual([batch[0].value for batch in self.sink.batches], [1.0, 2.0]) + + def test_flush_stops_at_the_first_failure(self): + """Leave the rest of the backlog queued when a retry fails again.""" + self.sink.failing = True + for index in range(3): + self.writer.write([_sample(value=float(index))]) + self.clock.advance(100.0) + self.assertEqual(self.writer.flush_due(), 0) + self.assertEqual(self.writer.queue_depth, 3) + + def test_reconnect_clears_backoff(self): + """Let an operator-driven reconnect retry immediately.""" + self.sink.failing = True + self.writer.write([_sample()]) + self.sink.failing = False + self.writer.connect() + self.assertEqual(self.sink.connects, 1) + self.assertEqual(self.writer.flush_due(), 1) + + def test_delegates_connection_state_and_close(self): + """Pass health and teardown through to the wrapped sink.""" + self.assertTrue(self.writer.is_connected()) + self.sink.failing = True + self.assertFalse(self.writer.is_connected()) + self.writer.close() + self.assertTrue(self.sink.closed) + + def test_policy_rejects_a_queue_that_holds_nothing(self): + """Refuse a queue bound that could never retain a batch.""" + with self.assertRaises(ValueError): + RetryPolicy(max_batches=0) + + def test_policy_rejects_a_non_positive_backoff(self): + """Refuse a backoff that would retry a dead backend without pause.""" + with self.assertRaises(ValueError): + RetryPolicy(base_backoff_s=0.0) + + +if __name__ == "__main__": + unittest.main() diff --git a/tox.ini b/tox.ini index 917c23a..6e2fa84 100644 --- a/tox.ini +++ b/tox.ini @@ -5,7 +5,8 @@ isolated_build = true [testenv] description = Run the test suite -extras = +# Installed so the InfluxDB sink's tests run rather than skip themselves +extras = influxdb deps = commands = python -m unittest discover -s tests