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
5 changes: 3 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,9 @@ Libby: a tiny messaging library which uses Bamboo with pluggable transports (ZMQ

## Documentation

Full docs (installation, keywords, the `Client` library, the `libby` CLI, and
how to build a `LibbyDaemon` peer, plus the generated API reference) are
Full docs (installation, keywords, the `Client` library, the `libby` CLI, how
to build a `LibbyDaemon` peer, the keygrabber, plus the generated API
reference) are
built with Sphinx + the Shibuya theme and published to GitHub Pages:

**[caltechopticalobservatories.github.io/libby](https://caltechopticalobservatories.github.io/libby/)**
Expand Down
18 changes: 18 additions & 0 deletions docs/source/api/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,24 @@ the docs build does not install the optional `influxdb` extra.
:show-inheritance:
```

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

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

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

## Responses and errors

```{eval-rst}
Expand Down
8 changes: 6 additions & 2 deletions docs/source/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,15 @@ Libby is a tiny messaging library built on [Bamboo](https://github.com/CaltechOp
with pluggable transports (ZMQ or RabbitMQ). It gives you:

- **Keywords** — typed, named values (`show` / `modify`) served over RPC, with
a registry, auto-generated `keys.list` / `keys.describe` services, and CLI
coercion.
a registry, auto-generated `keys.list` / `keys.describe` / `keys.read`
services, and CLI coercion.
- **`LibbyDaemon`** — a base class for peers: lifecycle, discovery, RPC
handlers, and pub/sub, in a few overrides.
- **`Client`** — a long-lived, in-process handle for reading and writing
keywords from scripts, and for blocking on a keyword condition with
`wait_for` (libby's `ktl.waitFor`).
- **Keygrabber**: a daemon that polls keywords from other peers and writes them
to a time-series database for dashboarding.
- **`libby` CLI** — a command-line front end for keyword peers
(`show` / `modify` / `list` / `describe` / `waitfor`).

Expand All @@ -23,6 +25,7 @@ keywords
client
cli
daemon
keygrabber
api/index
```

Expand All @@ -32,4 +35,5 @@ api/index
- Building a peer that serves keywords? Read {doc}`keywords` then {doc}`daemon`.
- Writing a script or tool that talks to peers? Read {doc}`client`.
- Poking at peers interactively? Read {doc}`cli`.
- Recording keywords for Grafana? Read {doc}`keygrabber`.
- Looking for a specific class or function? See {doc}`api/index`.
144 changes: 144 additions & 0 deletions docs/source/keygrabber.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
# Keygrabber

The keygrabber is a libby daemon that polls keywords from other peers on a
configurable cadence and writes them to a time-series database, so Grafana can
dashboard an instrument without every daemon growing its own database code.

One process serves the whole fleet. Cadence lives in this daemon's config
rather than in each hardware daemon, adding a metric never restarts a daemon
that owns moving hardware, and the database credential lives in one place.

InfluxDB 2.x is the only backend today. It needs the optional extra:

```bash
pip install libby[influxdb]
```

## Running it

```bash
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.

## Config

A keygrabber config is an ordinary daemon config (`peer_id`, `group_id`,
`transport`, and the transport's own settings) plus three sections of its own.

```yaml
peer_id: keygrabber
group_id: hispec
transport: rabbitmq
rabbitmq_url: amqp://localhost

sink:
type: influxdb
url: http://influx.hispec:8086
org: hispec
bucket: telemetry
token_env: HISPEC_INFLUX_TOKEN

workers: 4

defaults:
interval_s: 10.0
timeout_s: 2.0
refresh_s: 300.0

collections:
adc:
peer: hsfei.adc
interval_s: 5.0
keywords: ["positionvalue%", "ismoving", "isconnected"]
pressure:
peer: hsfei.atcpress
interval_s: 60.0
keywords: ["%"]
exclude: ["units_code"]
```

### sink

`type` selects the backend; `influxdb` is the only one implemented.

The token is **never** written in the config. `token_env` names an environment
variable to read it from, and a `token` key in the file is rejected outright.
`timeout_ms` is optional.

### collections

One entry per peer and cadence; several entries may target the same peer at
different cadences. A collection name must match `[a-z0-9_]+`, because it
becomes the prefix of that collection's control keywords.

- `peer` is `<group>.<daemon>`, the address of one peer
- `keywords` is a list of names or `%` patterns to record
- `exclude` removes names the includes matched
- `interval_s`, `timeout_s` and `refresh_s` fall back to `defaults`

`uptime` and `lasterror` are excluded by default: `uptime` changes every second
and says nothing a timestamp does not, and `lasterror` is null most of the
time. Naming either one in `keywords` explicitly opts it back in.

Patterns are resolved against the live peer at startup and again every
`refresh_s`, so keywords added by a restarted daemon get picked up without
restarting the keygrabber.

### Cadence and timeouts

Config load rejects an `interval_s` at or below `1.5 x timeout_s`. bamboo waits
`timeout_s` for the acknowledgement and a further `timeout_s / 2` for the
reply, so one unanswered read can occupy a worker for one and a half timeouts,
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.

## How it reads

A tick is one `keys.read` request per peer, not one per keyword. This matters
more than the smaller number suggests: a daemon dispatches requests inline on
its receive thread, so reading twenty keywords individually does not overlap
anything on that daemon. It serializes exactly as a batch would, while paying
twenty dispatch cycles instead of one. The saving is contention on a control
daemon's only dispatch thread.

A peer running a libby without `keys.read` is read one keyword at a time
instead. That is detected from the `services` field of `keys.list`, not by
trying `keys.read` and seeing what happens: an unknown key is dropped without
an acknowledgement, so a probe cannot tell an old peer from a dead one.

Every value in a tick carries one timestamp, taken by the keygrabber rather
than by each daemon, which keeps clock skew between daemon hosts out of the
data.

## Storage

Samples reach the backend through a `Sink`, and nothing in a `Sample` is
Influx-shaped, so a second backend is a new sink rather than a change to the
collector.

The InfluxDB schema is one measurement per keyword name, tagged with `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.

Three details follow from what Influx can store:

- Integers are written as floats, so one peer reporting `0` and another `0.5`
for the same keyword cannot collide as int against float and be rejected.
- A null value is skipped, since Influx has no null field. It is not an error.
- A keyword with no units gets a `units=none` tag, because Influx drops an
empty tag value and the keyword would otherwise split into two series.

A failed write goes to a bounded retry queue with exponential backoff, which
drops its oldest batch when full, so a database that stays down cannot grow the
daemon's memory without limit.

See {mod}`libby.keygrabber.sink` in the {doc}`API reference </api/index>` for
the sink contract.
17 changes: 17 additions & 0 deletions libby/keygrabber/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,15 @@
(``pip install libby[influxdb]``), so a deployment that only uses another
backend never has to install it.
"""
from .collection import Collection, TickResult
from .config import (
CollectionConfig,
KeygrabberConfig,
build_sink,
parse_config,
select_keywords,
)
from .daemon import KeygrabberDaemon
from .sink import (
RetryingWriter,
RetryPolicy,
Expand All @@ -15,10 +24,18 @@
)

__all__ = [
"Collection",
"CollectionConfig",
"KeygrabberConfig",
"KeygrabberDaemon",
"RetryingWriter",
"RetryPolicy",
"Sample",
"Sink",
"SinkError",
"SinkWriteError",
"TickResult",
"build_sink",
"parse_config",
"select_keywords",
]
145 changes: 145 additions & 0 deletions libby/keygrabber/collection.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
"""One configured peer, resolved against the live peer and read on a cadence."""
from __future__ import annotations

import time
from dataclasses import dataclass
from datetime import datetime
from typing import Any, Callable, Dict, List, Optional, Tuple

from ..client import Client
from ..errors import LibbyError, LibbyTimeout
from .config import CollectionConfig, select_keywords
from .sink import Sample

BULK_READ_SERVICE = "keys.read"


@dataclass(frozen=True)
class TickResult:
"""What one read of a collection produced."""

samples: Tuple[Sample, ...]
read_errors: int


class Collection:
"""Tracks what one peer exposes and turns a read of it into samples.

Resolution is refreshed periodically rather than once, so keywords added by
a restarted daemon are picked up without restarting the keygrabber.
"""

def __init__(
self,
config: CollectionConfig,
*,
clock: Callable[[], float] = time.monotonic,
) -> None:
self.config = config
self._clock = clock
self._names: Tuple[str, ...] = ()
self._bulk_read = False
self._resolved_at: Optional[float] = None

@property
def name(self) -> str:
"""Return the collection's configured name."""
return self.config.name

@property
def keyword_count(self) -> int:
"""Return how many keywords the last resolve selected."""
return len(self._names)

@property
def bulk_read(self) -> bool:
"""Return whether the peer advertised the bulk read service."""
return self._bulk_read

def needs_resolve(self) -> bool:
"""Return whether the keyword selection is due to be refreshed."""
if self._resolved_at is None:
return True
return self._clock() - self._resolved_at >= self.config.refresh_s

def resolve(self, client: Client) -> Tuple[str, ...]:
"""Ask the peer what it serves and select the configured keywords.

One ``keys.list`` covers both: the peer's keyword names, and whether it
serves ``keys.read``. Selection is then local, so a collection with
several patterns still costs one request.
"""
listing = client.listing(f"{self.config.peer}.%",
timeout_s=self.config.timeout_s)
available = [name.rsplit(".", 1)[-1] for name in listing.names]
self._names = select_keywords(self.config, available)
self._bulk_read = BULK_READ_SERVICE in listing.services
self._resolved_at = self._clock()
return self._names

def tick(self, client: Client, timestamp: datetime) -> TickResult:
"""Read every selected keyword once and return the samples."""
if not self._names:
return TickResult((), 0)

qualified = [f"{self.config.peer}.{name}" for name in self._names]
responses = (
client.read(qualified, timeout_s=self.config.timeout_s)
if self._bulk_read
else self._read_individually(client, qualified)
)

samples: List[Sample] = []
read_errors = 0
for qualified_name, response in responses.items():
if not response.get("ok"):
read_errors += 1
continue
samples.append(self._sample(qualified_name, response, timestamp))
return TickResult(tuple(samples), read_errors)

def _read_individually(
self,
client: Client,
qualified: List[str],
) -> Dict[str, Dict[str, Any]]:
"""Read one keyword at a time, for a peer without ``keys.read``.

Abandons the rest of the tick after the first timeout: a peer that has
stopped answering would otherwise cost ``timeout_s`` per keyword and
overrun the interval many times over.
"""
responses: Dict[str, Dict[str, Any]] = {}
timed_out = False
for name in qualified:
if timed_out:
responses[name] = {"ok": False, "error": "skipped after timeout"}
continue
try:
responses[name] = client.show(name, timeout_s=self.config.timeout_s)
except LibbyTimeout as exc:
timed_out = True
responses[name] = {"ok": False, "error": str(exc)}
except LibbyError as exc:
responses[name] = {"ok": False, "error": str(exc)}
return responses

def _sample(
self,
qualified_name: str,
response: Dict[str, Any],
timestamp: datetime,
) -> Sample:
"""Build a sample from one keyword response.

A null value is kept rather than dropped here: whether it can be stored
is the sink's business, and a backend other than Influx may hold it.
"""
return Sample(
keyword=qualified_name.rsplit(".", 1)[-1],
group=self.config.group,
peer=self.config.daemon,
value=response.get("value"),
units=response.get("units"),
timestamp=timestamp,
)
Loading
Loading