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
23 changes: 23 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,29 @@ See the [installation guide](docs/source/installation.md) for full setup
details, and the docs site above for everything else (keywords, the client
library, the CLI, and writing a `LibbyDaemon` peer).

## Finding peers

`libby list <group>.<daemon>` (and `Client.peers`) is how you find out which
peers are up. It broadcasts a `keys.list` that every peer answers, so it works
the same on both transports and needs nothing configured beyond the transport
itself.

Discovery, in the sense of a peer table you can ask who is alive, is not
implemented. It may be one day; until then use the broadcast above rather than
these:

- `Libby.peers_alive()` returns `{}` and `Libby.wait_for_peer()` always fails,
on either transport, because `Protocol` never instantiates a `PeerTable`.
- ZMQ runs bamboo's hello, but only toward peers already in the sender's
address book, and nothing consumes it for liveness. A client learns nothing
about a daemon that has not been told about the client, so
`Libby.knows_key()` stays False in the usual one-sided setup.
- `Libby.rabbitmq()` does not start hello at all. The broker routes messages;
it does not tell a peer who else is connected.

`LibbyDaemon`'s `discovery_enabled` / `on_hello` control bamboo's hello, not
the broadcast above.

## Testing

```bash
Expand Down
53 changes: 50 additions & 3 deletions docs/source/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,11 @@
```
libby show <group>.<daemon>.<keyword> # read a keyword (% wildcard in keyword)
libby modify <group>.<daemon>.<keyword>=V # write a keyword (exact keyword)
libby list <group>.<daemon>.<pattern> # list keyword names (% wildcard in keyword)
libby list <group>.<daemon>.<pattern> # list keyword names (% wildcard in any segment)
libby list <group>.<daemon> # list live daemons (% wildcard in any segment)
libby describe <group>.<daemon>.<keyword> # metadata for one keyword (exact keyword)
libby waitfor '$<group>.<daemon>.<keyword> > V' # block until a comparison holds
libby completion bash|zsh # print shell code for TAB completion
```

`<group>.<daemon>` is the address of one daemon: `group` is its `group_id`
Expand All @@ -16,8 +18,8 @@ daemon's config is addressed as `hsfei.adc`). `Libby.rabbitmq()` /
`Libby.zmq()` build the actual wire identity from those two fields via
`libby.naming.qualified_peer_id`, so a daemon config never needs to
concatenate them by hand. `group`/`daemon` are case-insensitive (`HSFEI.ADC`
and `hsfei.adc` reach the same daemon); the keyword isn't. Cross-peer
fanout is not supported. `req` and `sub` are kept for raw RPC / topic
and `hsfei.adc` reach the same daemon); the keyword isn't. `list` is the one
verb that spans daemons; `req` and `sub` are kept for raw RPC / topic
debugging.

## Examples
Expand Down Expand Up @@ -51,6 +53,16 @@ $ libby list hsfei.pickoff.%min
hsfei.pickoff.hardmin
hsfei.pickoff.softmin

$ libby list hsfei.% # which daemons are up?
hsfei.adc
hsfei.atcpress
hsfei.pickoff

$ libby list %.%.isconnected # one keyword across the whole fleet
hscal.hkettherm.isconnected
hsfei.adc.isconnected
hsfei.atcpress.isconnected

$ libby waitfor '$hsfei.pickoff.ismoving == false' --timeout 30
hsfei.pickoff.ismoving = False (satisfied after 4.2s)

Expand All @@ -62,6 +74,41 @@ Add `--json` to any verb for machine-readable output (objects for `show` /
`modify` / `describe`, list of objects for `show <pattern>`, list of strings
for `list`).

## Listing across daemons

A `%` in the `<group>` or `<daemon>` segment of `list` asks every reachable
daemon at once, rather than one named daemon. Drop the keyword segment
(`libby list hsfei.%`) to list the daemons themselves; keep it
(`libby list hsfei.%.is%`) to list matching keywords on each of them.

Nothing on the wire says how many daemons exist, so these always run for the
full timeout (default 1s, `--timeout` to change it) instead of returning on
the first answer. A daemon that is down simply doesn't appear. Exit code 3
means nothing answered.

Only daemons are listed. Every libby connection answers `keys.list`, so other
clients reply to the broadcast too, and they identify themselves as clients
and are left out.

Over ZMQ the broadcast only reaches daemons in the address book (`peers:` in
`cli_config.yaml`, or `--addr`); over RabbitMQ the broker reaches everyone.

## Completion

`libby completion bash` (or `zsh`) prints shell code that wires TAB
completion for verbs, flags and addresses. Add it to your shell rc:

```bash
eval "$(libby completion bash)"
```

Completing an address lists live daemons the same way `list` does, so TAB
offers daemons after `<group>.` and that daemon's keywords after
`<group>.<daemon>.`. Results are cached for 10s in
`~/.libby/completion_cache.json` so a burst of TABs costs one broadcast, and
the lookup is bounded at 0.5s so TAB never hangs. An unreachable broker
completes nothing rather than erroring.

## Modify syntax

- `key=value` or `key value` (positional) both work.
Expand Down
31 changes: 29 additions & 2 deletions docs/source/client.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,9 @@ client = Client.zmq(address_book={"hsfei_pickoff": "tcp://host:5555"})
units, flags); `set(name, value)` → the value the daemon applied.
- `wait_for(expression, timeout)` → blocks until a keyword satisfies a
comparison; see below.
- `list(pattern)` → matching qualified names; `describe(name)` → one keyword's
metadata; `read(names)` → many keywords in one request per peer; see below.
- `list(pattern)` → matching qualified names; `peers(pattern)` → the live
daemons; `describe(name)` → one keyword's metadata; `read(names)` → many
keywords in one request per peer; see below.
- Failures raise rather than return sentinels: `KeywordError` when the daemon
rejects a get/set (its message is on `.error`), `LibbyTimeout` when a
request isn't answered, both subclasses of `LibbyError`. `set` accepts
Expand Down Expand Up @@ -113,5 +114,31 @@ discovery registry, which stays empty without discovery, and calling
`keys.read` to see what happens cannot distinguish an old peer from a dead
one.

## Finding daemons

`peers` returns the live daemons matching a `<group>.<daemon>` pattern, and
`peer_listings` returns each one's keywords from the same round trip:

```python
client.peers("hsfei.%") # ["hsfei.adc", "hsfei.atcpress", ...]
client.peers() # every daemon, any group
client.peer_listings("hsfei.%") # {"hsfei.adc": ["isconnected", ...], ...}
```

`list` accepts the same wildcards in its group and daemon segments, so
`client.list("hsfei.%.isconnected")` reads across the group.

These broadcast a single `keys.list`, which every daemon answers. Nothing
says how many daemons exist, so they wait the full `timeout_s` (1s by
default) rather than returning on the first reply, and a daemon that is down
is simply absent. Over ZMQ the broadcast reaches only the daemons in the
address book; over RabbitMQ the broker reaches all of them.

Only daemons come back. Every `Libby` serves `keys.list`, so another
`Client` answers the broadcast as well, but `LibbyDaemon` is the only thing
that builds its `Libby` with `is_daemon=True`, and a reply counts only if it
reports that. A peer on a libby predating the flag is not listed, so the
daemons have to be running a libby this new to be found.

See {mod}`libby.client` in the {doc}`API reference </api/index>` for the
full method signatures.
7 changes: 4 additions & 3 deletions docs/source/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,15 +6,16 @@ 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` / `keys.read`
services, and CLI coercion.
- **`LibbyDaemon`** — a base class for peers: lifecycle, discovery, RPC
handlers, and pub/sub, in a few overrides.
- **`LibbyDaemon`** — a base class for peers: lifecycle, 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`).
(`show` / `modify` / `list` / `describe` / `waitfor`), with TAB completion
and wildcard listing to find which peers are up.

```{toctree}
:maxdepth: 2
Expand Down
3 changes: 2 additions & 1 deletion libby/__init__.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from bamboo.protocol import Protocol
from bamboo.builder import MessageBuilder
from bamboo.keys import KeyRegistry
from .libby import Libby
from .libby import BroadcastReply, Libby
from .keyword import (
Keyword,
BoolKeyword,
Expand All @@ -25,6 +25,7 @@

__all__ = [
"Libby",
"BroadcastReply",
"Client",
"KeyListing",
"WaitResult",
Expand Down
87 changes: 87 additions & 0 deletions libby/cli/completion.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
"""Shell completion for libby addresses, fed by one cached broadcast ``keys.list``.

The listing comes from the same broadcast ``libby list`` uses, not from
bamboo's hello/discovery, which never reports who is alive.
"""
from __future__ import annotations

import json
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Callable, Dict, List, Optional

DEFAULT_CACHE_PATH = Path.home() / ".libby" / "completion_cache.json"
CACHE_TTL_S = 10.0
COMPLETE_TIMEOUT_S = 0.5

Listings = Dict[str, List[str]]
ListingsSource = Callable[[float], Listings]


@dataclass(frozen=True)
class CompletionCache:
"""Last peer listing on disk, so a burst of TABs costs one broadcast."""

path: Path = DEFAULT_CACHE_PATH
ttl_s: float = CACHE_TTL_S

def load(self, connection: str) -> Optional[Listings]:
"""Return the cached listings for ``connection`` while still fresh."""
try:
raw = json.loads(self.path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
if not isinstance(raw, dict) or raw.get("connection") != connection:
return None
if time.time() - float(raw.get("ts", 0.0)) > self.ttl_s:
return None
listings = raw.get("listings")
return listings if isinstance(listings, dict) else None

def store(self, connection: str, listings: Listings) -> None:
"""Write the listings for ``connection``; an unwritable cache just means no cache."""
try:
self.path.parent.mkdir(parents=True, exist_ok=True)
self.path.write_text(
json.dumps({"connection": connection, "ts": time.time(), "listings": listings}),
encoding="utf-8",
)
except OSError:
pass


def cached_listings(connection: str, fetch: ListingsSource, cache: CompletionCache) -> Listings:
"""Return the listings from the cache, or fetch, store and return fresh ones."""
listings = cache.load(connection)
if listings is None:
listings = fetch(COMPLETE_TIMEOUT_S)
cache.store(connection, listings)
return listings


def peer_candidates(prefix: str, listings: Listings) -> List[str]:
"""Complete a partial ``<group>.<daemon>``."""
return sorted(peer for peer in listings if peer.startswith(prefix.lower()))


def address_candidates(prefix: str, listings: Listings) -> List[str]:
"""Complete a partial ``<group>.<daemon>.<keyword>``.

Until the daemon segment is complete the candidates end in ``.``, so the
shell stops there and the next TAB moves on to the keywords.
"""
if prefix.count(".") < 2:
peers = peer_candidates(prefix, listings)
# The shell appends a space to a lone candidate unless it ends in
# =/:, so an unambiguous daemon expands to its keywords right away
if len(peers) != 1:
return [f"{peer}." for peer in peers]
prefix = f"{peers[0]}."
peer, _, keyword_prefix = prefix.rpartition(".")
peer = peer.lower()
return sorted(
f"{peer}.{name}"
for name in listings.get(peer, [])
if name.startswith(keyword_prefix)
)
Loading
Loading