From 4d4bf350018c848145ea98c4c04950596b225829 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:16:26 -0700 Subject: [PATCH 1/6] Add wildcard peer listing and shell completion to the libby CLI --- docs/source/cli.md | 49 +++++++++++- docs/source/client.md | 25 +++++- libby/__init__.py | 3 +- libby/cli/completion.py | 83 +++++++++++++++++++ libby/cli/libby_cli.py | 133 ++++++++++++++++++++++++++----- libby/client.py | 55 +++++++++++-- libby/libby.py | 59 +++++++++++++- libby/naming.py | 35 ++++++++ pyproject.toml | 1 + tests/test_client_integration.py | 125 ++++++++++++++++++++++++++++- tests/test_completion.py | 97 ++++++++++++++++++++++ tests/test_naming.py | 30 ++++++- 12 files changed, 660 insertions(+), 35 deletions(-) create mode 100644 libby/cli/completion.py create mode 100644 tests/test_completion.py diff --git a/docs/source/cli.md b/docs/source/cli.md index 013e0e5..8bf2de0 100644 --- a/docs/source/cli.md +++ b/docs/source/cli.md @@ -5,9 +5,11 @@ ``` libby show .. # read a keyword (% wildcard in keyword) libby modify ..=V # write a keyword (exact keyword) -libby list .. # list keyword names (% wildcard in keyword) +libby list .. # list keyword names (% wildcard in any segment) +libby list . # list live daemons (% wildcard in any segment) libby describe .. # metadata for one keyword (exact keyword) libby waitfor '$.. > V' # block until a comparison holds +libby completion bash|zsh # print shell code for TAB completion ``` `.` is the address of one daemon: `group` is its `group_id` @@ -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 @@ -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) @@ -62,6 +74,37 @@ Add `--json` to any verb for machine-readable output (objects for `show` / `modify` / `describe`, list of objects for `show `, list of strings for `list`). +## Listing across daemons + +A `%` in the `` or `` 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. + +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 discovers live daemons the same way `list` does, so +TAB offers daemons after `.` and that daemon's keywords after +`..`. Results are cached for 10s in +`~/.libby/completion_cache.json` so a burst of TABs costs one broadcast, and +discovery 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. diff --git a/docs/source/client.md b/docs/source/client.md index 614996c..b1f0b44 100644 --- a/docs/source/client.md +++ b/docs/source/client.md @@ -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 @@ -113,5 +114,25 @@ 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 `.` 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. + See {mod}`libby.client` in the {doc}`API reference ` for the full method signatures. diff --git a/libby/__init__.py b/libby/__init__.py index 139d2b3..f3da893 100644 --- a/libby/__init__.py +++ b/libby/__init__.py @@ -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, @@ -25,6 +25,7 @@ __all__ = [ "Libby", + "BroadcastReply", "Client", "KeyListing", "WaitResult", diff --git a/libby/cli/completion.py b/libby/cli/completion.py new file mode 100644 index 0000000..cea5bca --- /dev/null +++ b/libby/cli/completion.py @@ -0,0 +1,83 @@ +"""Shell completion for libby addresses, fed by one cached broadcast ``keys.list``.""" +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 discovery result 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 ``.``.""" + return sorted(peer for peer in listings if peer.startswith(prefix.lower())) + + +def address_candidates(prefix: str, listings: Listings) -> List[str]: + """Complete a partial ``..``. + + 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) + ) diff --git a/libby/cli/libby_cli.py b/libby/cli/libby_cli.py index 01041d7..749d3e6 100644 --- a/libby/cli/libby_cli.py +++ b/libby/cli/libby_cli.py @@ -1,4 +1,5 @@ """Libby CLI — show/modify keywords on libby peers, plus raw req/sub.""" +# PYTHON_ARGCOMPLETE_OK from __future__ import annotations import argparse @@ -10,9 +11,19 @@ import sys import time from pathlib import Path -from typing import Any, Dict, List, Optional, Tuple +from typing import Any, Callable, Dict, List, Optional, Tuple from urllib.parse import urlsplit, urlunsplit +import argcomplete +from argcomplete.shell_integration import shellcode + +from libby.cli.completion import ( + CompletionCache, + Listings, + address_candidates, + cached_listings, + peer_candidates, +) from libby.config_resolve import ( DEFAULT_BIND, DEFAULT_CONFIG_PATH, @@ -22,11 +33,11 @@ resolve_rabbitmq_url, resolve_transport, ) -from libby.client import DEFAULT_POLL_S, Client, WaitResult +from libby.client import DEFAULT_BROADCAST_TIMEOUT_S, DEFAULT_POLL_S, Client, WaitResult from libby.errors import KeywordError, LibbyError from libby.expression import parse_comparison from libby.libby import Libby -from libby.naming import coerce_value, parse_keyword, peer_id +from libby.naming import coerce_value, parse_address_pattern, parse_keyword, peer_id from libby.response import unwrap DEFAULT_SELF_ID = "cli" @@ -347,15 +358,22 @@ def cmd_show(namespace: argparse.Namespace) -> int: def cmd_list(namespace: argparse.Namespace) -> int: config = load_cli_config(namespace.config) - # Reject a malformed address here, before opening a transport, so it stays - # an argument error; Client.list parses it again for the peer id - parse_keyword(namespace.pattern, allow_pattern=True) - timeout = namespace.timeout if namespace.timeout is not None else DEFAULT_TIMEOUT_S + # Parse before opening a transport so a malformed pattern stays an argument error + address = parse_address_pattern(namespace.pattern) + # A pattern that spans daemons is answered by broadcast, which always runs + # to its timeout, so it gets the shorter discovery default + one_daemon = address.keyword is not None and not address.spans_peers + default_timeout = DEFAULT_TIMEOUT_S if one_daemon else DEFAULT_BROADCAST_TIMEOUT_S + timeout = namespace.timeout if namespace.timeout is not None else default_timeout lib: Optional[Libby] = None try: lib = _mk_libby(namespace, config) - qualified_names = Client(lib).list(namespace.pattern, timeout_s=timeout) + client = Client(lib) + if address.keyword is None: + qualified_names = client.peers(namespace.pattern, timeout_s=timeout) + else: + qualified_names = client.list(namespace.pattern, timeout_s=timeout) if not qualified_names: if namespace.json: print(json.dumps([], indent=2)) @@ -570,6 +588,71 @@ def _printer(msg): print("[libby sub] stopped") +def _connection_key(namespace: argparse.Namespace, config: Dict[str, Any]) -> str: + """Name the broker or address book a completion cache entry belongs to.""" + if resolve_transport(namespace.transport, config) == "rabbitmq": + return f"rabbitmq {_redact_url(resolve_rabbitmq_url(namespace.rabbitmq_url, config))}" + return f"zmq {sorted(resolve_address_book(config, namespace.addr).items())}" + + +def _discover_listings(namespace: argparse.Namespace) -> Listings: + config = load_cli_config(namespace.config) + + def fetch(timeout_s: float) -> Listings: + lib = _mk_libby(namespace, config) + try: + return Client(lib).peer_listings(timeout_s=timeout_s) + finally: + lib.stop() + + return cached_listings(_connection_key(namespace, config), fetch, CompletionCache()) + + +# argcomplete also passes ``action`` and ``parser``, hence the catch-all +def _complete_address(prefix: str, parsed_args: argparse.Namespace, **_: Any) -> List[str]: + """Complete a partial .. argument.""" + try: + return address_candidates(prefix, _discover_listings(parsed_args)) + except Exception: # pylint: disable=broad-exception-caught + # A completer must never break the shell; an unreachable broker completes nothing + return [] + + +def _complete_peer(prefix: str, parsed_args: argparse.Namespace, **_: Any) -> List[str]: + """Complete a partial . argument.""" + try: + return peer_candidates(prefix, _discover_listings(parsed_args)) + except Exception: # pylint: disable=broad-exception-caught + return [] + + +def _completed_by( + action: argparse.Action, + completer: Callable[..., List[str]] = _complete_address, +) -> argparse.Action: + """Attach a completer to an argument; argcomplete finds it by attribute.""" + setattr(action, "completer", completer) + return action + + +def cmd_completion(namespace: argparse.Namespace) -> int: + """Print the shell code that wires TAB completion to the libby executable.""" + print(shellcode(["libby"], shell=namespace.shell)) + return 0 + + +def _add_completion_verb(sub: argparse._SubParsersAction) -> None: + """Register the completion verb, the one verb with no connection flags.""" + parser = sub.add_parser( + "completion", + help="Print shell code that enables TAB completion; eval it from your shell rc", + ) + parser.add_argument("shell", choices=("bash", "zsh")) + # main() reads the common logging flags, which this verb has no use for + parser.set_defaults(func=cmd_completion, config=None, log_level=None, + log_file=None, verbose=False) + + def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser( prog="libby", @@ -607,28 +690,31 @@ def add_common(p): # argparse %-formats help strings, so a literal % has to be escaped p_show = sub.add_parser("show", help="Read a keyword's value (%% allowed in keyword)") add_common(p_show) - p_show.add_argument("keyword", - help=".. (%% allowed in keyword segment)") + _completed_by(p_show.add_argument( + "keyword", help=".. (%% allowed in keyword segment)")) p_show.set_defaults(func=cmd_show) - p_list = sub.add_parser("list", help="List keyword names matching a pattern") + p_list = sub.add_parser("list", help="List keyword names, or daemons, matching a pattern") add_common(p_list) - p_list.add_argument("pattern", - help=".. (%% wildcard in keyword)") + _completed_by(p_list.add_argument( + "pattern", + help=".[.] (%% wildcard in any segment; " + "without a keyword segment, list the matching daemons)", + )) p_list.set_defaults(func=cmd_list) p_describe = sub.add_parser("describe", help="Show metadata for a keyword") add_common(p_describe) - p_describe.add_argument("keyword", - help=".. (exact, no wildcards)") + _completed_by(p_describe.add_argument( + "keyword", help=".. (exact, no wildcards)")) p_describe.set_defaults(func=cmd_describe) p_modify = sub.add_parser("modify", help="Set a keyword's value") add_common(p_modify) - p_modify.add_argument( + _completed_by(p_modify.add_argument( "keyword", help="..= or .. (+ value arg)", - ) + )) p_modify.add_argument("value", nargs="?", help="Value (if not using = form)") p_modify.set_defaults(func=cmd_modify) @@ -649,9 +735,10 @@ def add_common(p): help="'$.. '; " "op is ==, !=, <, <=, >, >=", ) - p_waitfor.add_argument("-d", "--daemon", metavar=".", - help="Default daemon, so the expression can name a " - "keyword bare (e.g. '$ismoving == false')") + _completed_by(p_waitfor.add_argument( + "-d", "--daemon", metavar=".", + help="Default daemon, so the expression can name a " + "keyword bare (e.g. '$ismoving == false')"), _complete_peer) p_waitfor.add_argument("--case", action="store_true", help="Compare strings case-sensitively " "(default: case-insensitive)") @@ -672,11 +759,15 @@ def add_common(p): p_sub.add_argument("topics", nargs="+", help="Topic(s) to subscribe to") p_sub.set_defaults(func=cmd_sub) + _add_completion_verb(sub) + return parser def main(argv: Optional[List[str]] = None) -> int: - namespace = build_parser().parse_args(argv) + parser = build_parser() + argcomplete.autocomplete(parser) + namespace = parser.parse_args(argv) try: cli_config = load_cli_config(namespace.config) diff --git a/libby/client.py b/libby/client.py index 3a15a78..a85910e 100644 --- a/libby/client.py +++ b/libby/client.py @@ -18,10 +18,11 @@ resolve_rabbitmq_url, resolve_transport, ) -from .errors import LibbyError, LibbyTimeout +from .errors import KeywordNameError, LibbyError, LibbyTimeout from .expression import parse_comparison -from .libby import Libby -from .naming import parse_keyword, peer_id +from .keyword import match_pattern +from .libby import DEFAULT_BROADCAST_TIMEOUT_S, Libby +from .naming import parse_address_pattern, parse_keyword, peer_id from .response import unwrap DEFAULT_SELF_ID = "libby-client" @@ -163,9 +164,53 @@ def list(self, pattern: str, *, timeout_s: float = DEFAULT_TIMEOUT_S) -> List[st """List qualified keyword names matching ``..``. Returns fully qualified names, so the result feeds straight back into - :meth:`get`, :meth:`show` or :meth:`read`. + :meth:`get`, :meth:`show` or :meth:`read`. A ``%`` in the group or + daemon segment spans daemons by broadcast, so that form waits the + full ``timeout_s`` (see :meth:`peers`). """ - return list(self.listing(pattern, timeout_s=timeout_s).names) + address = parse_address_pattern(pattern) + if address.keyword is None: + raise KeywordNameError(f"list needs a keyword pattern: {pattern}") + if not address.spans_peers: + return list(self.listing(pattern, timeout_s=timeout_s).names) + listings = self.peer_listings(pattern, timeout_s=timeout_s) + return [f"{peer}.{name}" for peer in sorted(listings) for name in listings[peer]] + + def peers( + self, + pattern: str = "%.%", + *, + timeout_s: float = DEFAULT_BROADCAST_TIMEOUT_S, + ) -> List[str]: + """List the live daemons matching ``.`` as ``group.daemon`` ids. + + Every daemon answers a broadcast ``keys.list``, so this waits the full + ``timeout_s`` rather than returning on the first reply. Over ZMQ only + daemons in the address book are asked. + """ + return sorted(self.peer_listings(pattern, timeout_s=timeout_s)) + + def peer_listings( + self, + pattern: str = "%.%", + *, + timeout_s: float = DEFAULT_BROADCAST_TIMEOUT_S, + ) -> Dict[str, List[str]]: + """Map each live daemon matching ``pattern`` to its keyword names. + + A keyword segment in ``pattern`` narrows the names; without one every + keyword is listed. Same timeout semantics as :meth:`peers`. + """ + address = parse_address_pattern(pattern) + replies = self._libby.broadcast_request( + "keys.list", {"pattern": address.keyword or "%"}, timeout_s=timeout_s) + wanted = set(match_pattern(address.peer_pattern, + (reply.peer_id for reply in replies))) + return { + reply.peer_id: list(reply.payload.get("matches", [])) + for reply in replies + if reply.peer_id in wanted and reply.payload.get("ok") + } def describe(self, name: str, *, timeout_s: float = DEFAULT_TIMEOUT_S) -> Dict[str, Any]: """Read one keyword's metadata (type, access, units, timeout_s).""" diff --git a/libby/libby.py b/libby/libby.py index c946d80..695480d 100644 --- a/libby/libby.py +++ b/libby/libby.py @@ -1,7 +1,11 @@ from __future__ import annotations -from typing import Any, Callable, Dict, Iterable, List, Optional +from contextlib import contextmanager +from dataclasses import dataclass +from typing import Any, Callable, Dict, Iterable, Iterator, List, Optional +import queue import time +from bamboo.builder import MessageBuilder from bamboo.keys import KeyRegistry from bamboo.protocol import Protocol from bamboo.discovery import Discovery @@ -10,6 +14,17 @@ from .keyword_registry import KeywordRegistry from .naming import qualified_peer_id +DEFAULT_BROADCAST_TIMEOUT_S = 1.0 + + +@dataclass(frozen=True) +class BroadcastReply: + """One peer's answer to :meth:`Libby.broadcast_request`.""" + + peer_id: str + payload: Dict[str, Any] + + class Libby: def __init__( self, @@ -167,6 +182,48 @@ def request(self, peer_id: str, key: str, payload: Dict[str, Any], ttl_ms: int = def rpc(self, peer_id: str, key: str, payload: Dict[str, Any], ttl_ms: int = 8000): return self.request(peer_id, key, payload, ttl_ms) + def broadcast_request( + self, + key: str, + payload: Dict[str, Any], + timeout_s: float = DEFAULT_BROADCAST_TIMEOUT_S, + ) -> List[BroadcastReply]: + """Send ``key`` to every reachable peer and collect their replies. + + Nothing on the wire says how many peers will answer, so this always + waits the full ``timeout_s``. Over ZMQ "every reachable peer" is the + address book. Our own reply is dropped: on RabbitMQ our queue sits on + the fanout exchange like everyone else's. + """ + msg = MessageBuilder(sourceid=self.self_id).req(key, payload).to(None).build() + with self._reply_queue(msg.env.transid) as replies: + self.proto.send(msg) + return list(self._drain_replies(replies, timeout_s)) + + @contextmanager + def _reply_queue(self, transid: str) -> Iterator["queue.Queue[Any]"]: + # bamboo only registers a reply queue for direct requests, so a broadcast needs its own + lock, waiters = self.proto._lock, self.proto._resp_wait # pylint: disable=protected-access + replies: "queue.Queue[Any]" = queue.Queue() + with lock: + waiters[transid] = replies + try: + yield replies + finally: + with lock: + waiters.pop(transid, None) + + def _drain_replies(self, replies: "queue.Queue[Any]", + timeout_s: float) -> Iterator[BroadcastReply]: + deadline = time.monotonic() + timeout_s + while (remaining := deadline - time.monotonic()) > 0: + try: + resp = replies.get(timeout=remaining) + except queue.Empty: + return + if resp.env.sourceid != self.self_id: + yield BroadcastReply(resp.env.sourceid, dict(resp.env.payload or {})) + def serve_keys(self, keys: List[str], callback: Callable[[dict, dict], Optional[dict]]) -> None: for k in keys: self.proto.serve(k, callback) diff --git a/libby/naming.py b/libby/naming.py index 06da89b..6ef7395 100644 --- a/libby/naming.py +++ b/libby/naming.py @@ -1,6 +1,7 @@ """Keyword-name parsing, value coercion, and peer/group naming (transport-agnostic).""" from __future__ import annotations +from dataclasses import dataclass from typing import Any, Optional, Tuple from .errors import KeywordNameError @@ -29,6 +30,40 @@ def parse_keyword(arg: str, *, allow_pattern: bool = False) -> Tuple[str, str, s return group, daemon, keyword +@dataclass(frozen=True) +class AddressPattern: + """A ``.[.]`` pattern with ``%`` allowed in every segment.""" + + group: str + daemon: str + keyword: Optional[str] + + @property + def spans_peers(self) -> bool: + """Return True when the group or daemon segment holds a wildcard.""" + return "%" in self.group or "%" in self.daemon + + @property + def peer_pattern(self) -> str: + """Return the pattern to match wire ids against, lowercased like the ids.""" + return f"{self.group}.{self.daemon}".lower() + + +def parse_address_pattern(arg: str) -> AddressPattern: + """Parse ``.[.]`` with ``%`` allowed anywhere. + + Unlike :func:`parse_keyword` this admits patterns that span daemons, so it + is only for callers that resolve them by broadcast (``list``, ``peers``). + """ + parts = arg.split(".", 2) + if len(parts) < 2 or not all(parts): + raise KeywordNameError( + f"pattern must be .[.], got: {arg}" + ) + keyword = parts[2] if len(parts) == 3 else None + return AddressPattern(parts[0], parts[1], keyword) + + def qualified_peer_id(peer_id: str, group_id: Optional[str] = None) -> str: """Return the lowercased wire identity ".", or bare peer_id. diff --git a/pyproject.toml b/pyproject.toml index 87ebed0..17cc390 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -18,6 +18,7 @@ dependencies = [ "pyzmq>=25.0.0", "pika>=1.3.0", "pyyaml>=6.0", + "argcomplete>=3.0", "bamboo @ git+https://github.com/CaltechOpticalObservatories/bamboo.git@main" ] diff --git a/tests/test_client_integration.py b/tests/test_client_integration.py index 157d742..ee304d5 100644 --- a/tests/test_client_integration.py +++ b/tests/test_client_integration.py @@ -14,15 +14,18 @@ import socket import threading import unittest -from typing import Tuple +from typing import List, Optional, Tuple, Type from libby import Client, KeywordError +from libby.client import DEFAULT_SELF_ID from libby.daemon import LibbyDaemon from libby.rabbitmq_transport import RabbitMQTransport RABBITMQ_URL = "amqp://localhost" GROUP_ID = "hsfei" +OTHER_GROUP_ID = "hscal" RPC_TIMEOUT_S = 6.0 +DISCOVERY_TIMEOUT_S = 2.0 SERVE_STOP_TIMEOUT_S = 10.0 INITIAL_POSITION = 10.0 @@ -104,6 +107,22 @@ class _ZmqLifecycleDaemon(_LifecycleDaemon): transport = "zmq" +def _started_daemon( + daemon_cls: Type[_FixtureDaemon], + daemon_peer_id: str, + daemon_group_id: str, + bind: Optional[str] = None, +) -> LibbyDaemon: + """Start one fixture daemon under its own peer and group id.""" + daemon = daemon_cls() + daemon.peer_id = daemon_peer_id + daemon.group_id = daemon_group_id + if bind is not None: + daemon.bind = bind + daemon.start() + return daemon + + # Nested inside a plain class so unittest's loader, which collects every # module-level TestCase subclass, does not run the bases on their own. class _Bases: # pylint: disable=too-few-public-methods @@ -206,6 +225,61 @@ def test_read_merges_chunked_requests(self): self.assertEqual(list(values), wanted) self.assertTrue(all(entry["ok"] for entry in values.values())) + class DiscoveryCases(unittest.TestCase): + """Broadcast discovery that must hold identically on every transport. + + Concrete subclasses start two fixture daemons in ``GROUP_ID`` and one + in ``OTHER_GROUP_ID``, and supply ``client`` plus their qualified ids. + """ + + client: Client + group_peers: Tuple[str, str] + other_peer: str + + def test_peers_lists_every_daemon_in_the_group(self): + """Find the group's daemons by wildcard, and no other group's.""" + found = self.client.peers(f"{GROUP_ID}.%", timeout_s=DISCOVERY_TIMEOUT_S) + for peer in self.group_peers: + self.assertIn(peer, found) + self.assertNotIn(self.other_peer, found) + + def test_peers_spans_groups_and_skips_the_client(self): + """Find every daemon with %.% without listing the asking client.""" + found = self.client.peers("%.%", timeout_s=DISCOVERY_TIMEOUT_S) + for peer in (*self.group_peers, self.other_peer): + self.assertIn(peer, found) + self.assertNotIn(DEFAULT_SELF_ID, found) + + def test_peers_with_an_exact_id_confirms_one_daemon(self): + """Resolve an exact . to just that daemon.""" + found = self.client.peers(self.group_peers[0], timeout_s=DISCOVERY_TIMEOUT_S) + self.assertEqual(found, [self.group_peers[0]]) + + def test_list_across_daemons_returns_qualified_names(self): + """List one keyword on every daemon of a group as names that feed back in.""" + names = self.client.list(f"{GROUP_ID}.%.positionvalue", + timeout_s=DISCOVERY_TIMEOUT_S) + for peer in self.group_peers: + self.assertIn(f"{peer}.positionvalue", names) + self.assertNotIn(f"{self.other_peer}.positionvalue", names) + self.assertEqual( + self.client.get(f"{self.group_peers[0]}.positionvalue", + timeout_s=RPC_TIMEOUT_S), + INITIAL_POSITION) + + def test_list_across_daemons_with_no_match_is_empty(self): + """Return nothing, rather than raise, when no daemon has the keyword.""" + self.assertEqual( + self.client.list(f"{GROUP_ID}.%.nosuchkeyword", timeout_s=DISCOVERY_TIMEOUT_S), + []) + + def test_peer_listings_carry_each_daemons_keywords(self): + """Map each daemon to its keyword names from the same broadcast.""" + listings = self.client.peer_listings(f"{GROUP_ID}.%", timeout_s=DISCOVERY_TIMEOUT_S) + for peer in self.group_peers: + self.assertIn("positionvalue", listings[peer]) + self.assertIn("uptime", listings[peer]) + class LifecycleCases(unittest.TestCase): """Daemon startup and shutdown guarantees, per transport.""" @@ -283,6 +357,55 @@ def tearDownClass(cls): cls.daemon.stop() +@unittest.skipUnless(_broker_available(), "no RabbitMQ broker reachable at amqp://localhost") +class RabbitMQDiscoveryTests(_Bases.DiscoveryCases): + """Discovery cases over RabbitMQ.""" + + daemons: List[LibbyDaemon] + + @classmethod + def setUpClass(cls): + cls.daemons = [ + _started_daemon(_RabbitFixtureDaemon, "discoveryonermq", GROUP_ID), + _started_daemon(_RabbitFixtureDaemon, "discoverytwormq", GROUP_ID), + _started_daemon(_RabbitFixtureDaemon, "discoveryotherrmq", OTHER_GROUP_ID), + ] + cls.group_peers = (f"{GROUP_ID}.discoveryonermq", f"{GROUP_ID}.discoverytwormq") + cls.other_peer = f"{OTHER_GROUP_ID}.discoveryotherrmq" + cls.client = Client.rabbitmq(rabbitmq_url=RABBITMQ_URL) + + @classmethod + def tearDownClass(cls): + cls.client.close() + for daemon in cls.daemons: + daemon.stop() + + +class ZmqDiscoveryTests(_Bases.DiscoveryCases): + """Discovery cases over ZMQ, where the address book is what gets asked.""" + + daemons: List[LibbyDaemon] + + @classmethod + def setUpClass(cls): + groups = {"discoveryonezmq": GROUP_ID, "discoverytwozmq": GROUP_ID, + "discoveryotherzmq": OTHER_GROUP_ID} + address_book = {f"{group}.{name}": _free_endpoint() for name, group in groups.items()} + cls.daemons = [ + _started_daemon(_ZmqFixtureDaemon, name, group, bind=address_book[f"{group}.{name}"]) + for name, group in groups.items() + ] + cls.group_peers = (f"{GROUP_ID}.discoveryonezmq", f"{GROUP_ID}.discoverytwozmq") + cls.other_peer = f"{OTHER_GROUP_ID}.discoveryotherzmq" + cls.client = Client.zmq(bind=_free_endpoint(), address_book=address_book) + + @classmethod + def tearDownClass(cls): + cls.client.close() + for daemon in cls.daemons: + daemon.stop() + + @unittest.skipUnless(_broker_available(), "no RabbitMQ broker reachable at amqp://localhost") class RabbitMQLifecycleTests(_Bases.LifecycleCases): """Lifecycle cases over RabbitMQ.""" diff --git a/tests/test_completion.py b/tests/test_completion.py new file mode 100644 index 0000000..0a715ee --- /dev/null +++ b/tests/test_completion.py @@ -0,0 +1,97 @@ +"""Unit tests for the CLI's address completion and its cache. No transport needed.""" +import tempfile +import time +import unittest +from pathlib import Path + +from libby.cli.completion import ( + CompletionCache, + address_candidates, + cached_listings, + peer_candidates, +) + +LISTINGS = { + "hsfei.adc": ["isconnected", "position"], + "hsfei.atcfw": ["isconnected"], + "hscal.hkettherm": ["t2A_red_etalon"], +} +CONNECTION = "rabbitmq amqp://localhost" + + +class AddressCandidatesTests(unittest.TestCase): + def test_partial_group_offers_daemons_with_a_trailing_dot(self): + self.assertEqual(address_candidates("hs", LISTINGS), + ["hscal.hkettherm.", "hsfei.adc.", "hsfei.atcfw."]) + + def test_unambiguous_daemon_expands_to_its_keywords(self): + # Never a lone ".." candidate, which the shell would + # follow with a space + self.assertEqual(address_candidates("hsfei.at", LISTINGS), + ["hsfei.atcfw.isconnected"]) + + def test_complete_daemon_offers_its_keywords(self): + self.assertEqual(address_candidates("hsfei.adc.", LISTINGS), + ["hsfei.adc.isconnected", "hsfei.adc.position"]) + + def test_daemon_is_case_insensitive_and_keyword_is_not(self): + self.assertEqual(address_candidates("HSFEI.adc.pos", LISTINGS), ["hsfei.adc.position"]) + self.assertEqual(address_candidates("hscal.hkettherm.T", LISTINGS), []) + + def test_unknown_daemon_offers_nothing(self): + self.assertEqual(address_candidates("hsfei.nosuch.", LISTINGS), []) + + +class PeerCandidatesTests(unittest.TestCase): + def test_offers_matching_daemons_without_a_trailing_dot(self): + self.assertEqual(peer_candidates("hsfei.", LISTINGS), ["hsfei.adc", "hsfei.atcfw"]) + + def test_empty_prefix_offers_everything(self): + self.assertEqual(len(peer_candidates("", LISTINGS)), len(LISTINGS)) + + +class CompletionCacheTests(unittest.TestCase): + def setUp(self): + self._tmp = tempfile.TemporaryDirectory() # pylint: disable=consider-using-with + self.addCleanup(self._tmp.cleanup) + self.path = Path(self._tmp.name) / "nested" / "cache.json" + + def test_missing_file_is_a_miss(self): + self.assertIsNone(CompletionCache(self.path).load(CONNECTION)) + + def test_store_then_load_round_trips_and_creates_the_directory(self): + cache = CompletionCache(self.path) + cache.store(CONNECTION, LISTINGS) + self.assertEqual(cache.load(CONNECTION), LISTINGS) + + def test_another_connection_is_a_miss(self): + cache = CompletionCache(self.path) + cache.store(CONNECTION, LISTINGS) + self.assertIsNone(cache.load("zmq []")) + + def test_expired_entry_is_a_miss(self): + cache = CompletionCache(self.path, ttl_s=0.0) + cache.store(CONNECTION, LISTINGS) + time.sleep(0.01) + self.assertIsNone(cache.load(CONNECTION)) + + def test_corrupt_file_is_a_miss(self): + self.path.parent.mkdir(parents=True) + self.path.write_text("not json", encoding="utf-8") + self.assertIsNone(CompletionCache(self.path).load(CONNECTION)) + + def test_cached_listings_fetches_once_within_the_ttl(self): + calls = [] + + def fetch(timeout_s: float): + calls.append(timeout_s) + return LISTINGS + + cache = CompletionCache(self.path) + self.assertEqual(cached_listings(CONNECTION, fetch, cache), LISTINGS) + self.assertEqual(cached_listings(CONNECTION, fetch, cache), LISTINGS) + self.assertEqual(len(calls), 1) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_naming.py b/tests/test_naming.py index 07cff6d..9a0d335 100644 --- a/tests/test_naming.py +++ b/tests/test_naming.py @@ -2,7 +2,13 @@ import unittest from libby.errors import KeywordNameError -from libby.naming import coerce_value, parse_keyword, peer_id, qualified_peer_id +from libby.naming import ( + coerce_value, + parse_address_pattern, + parse_keyword, + peer_id, + qualified_peer_id, +) class ParseKeywordTests(unittest.TestCase): @@ -35,6 +41,28 @@ def test_rejects_wildcard_in_keyword_unless_allowed(self): ) +class ParseAddressPatternTests(unittest.TestCase): + def test_two_segments_mean_no_keyword(self): + address = parse_address_pattern("hsfei.%") + self.assertEqual((address.group, address.daemon, address.keyword), ("hsfei", "%", None)) + + def test_wildcard_daemon_spans_peers(self): + address = parse_address_pattern("HSFEI.%.is%") + self.assertTrue(address.spans_peers) + self.assertEqual(address.keyword, "is%") + # Wire ids are lowercased, so the pattern matched against them is too + self.assertEqual(address.peer_pattern, "hsfei.%") + + def test_exact_daemon_does_not_span_peers(self): + self.assertFalse(parse_address_pattern("hsfei.pickoff.is%").spans_peers) + + def test_rejects_one_segment_and_empty_segments(self): + with self.assertRaises(KeywordNameError): + parse_address_pattern("hsfei") + with self.assertRaises(KeywordNameError): + parse_address_pattern("hsfei..positionvalue") + + class PeerIdTests(unittest.TestCase): def test_joins_group_and_daemon_with_a_dot(self): # Must match the wire identity qualified_peer_id builds daemon-side From 1903658727f25662496007c1d47ec0fee889a260 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:40:52 -0700 Subject: [PATCH 2/6] Document that finding peers is the list broadcast, not bamboo discovery --- README.md | 21 +++++++++++++++++++++ docs/source/index.md | 7 ++++--- 2 files changed, 25 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index edb2e83..72c7ba1 100644 --- a/README.md +++ b/README.md @@ -32,6 +32,27 @@ 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 .` (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. + +Bamboo's own hello/discovery is a separate mechanism and is not a way to find +peers today: + +- `Protocol` never instantiates a `PeerTable`, so `Libby.peers_alive()` returns + `{}` and `Libby.wait_for_peer()` always fails, on either transport. +- `Libby.rabbitmq()` doesn't start discovery at all. The broker routes messages; + it does not tell a peer who else is connected. +- On ZMQ a hello only reaches peers already in the sender's address book, so a + client learns nothing about a daemon that hasn't been told about the client. + `Libby.knows_key()` stays False in the usual one-sided setup. + +`LibbyDaemon`'s `discovery_enabled` / `on_hello` control that mechanism, not +the broadcast above. + ## Testing ```bash diff --git a/docs/source/index.md b/docs/source/index.md index 4bb1a2d..1d0f0fd 100644 --- a/docs/source/index.md +++ b/docs/source/index.md @@ -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 From 0e5b5005bcac69aaefee5665a688aedf6839adc2 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:47:58 -0700 Subject: [PATCH 3/6] Stop calling the list broadcast discovery, which means bamboo's hello --- docs/source/cli.md | 6 ++--- libby/cli/completion.py | 8 ++++-- libby/cli/libby_cli.py | 8 +++--- tests/test_client_integration.py | 44 ++++++++++++++++---------------- 4 files changed, 35 insertions(+), 31 deletions(-) diff --git a/docs/source/cli.md b/docs/source/cli.md index 8bf2de0..79750b6 100644 --- a/docs/source/cli.md +++ b/docs/source/cli.md @@ -98,11 +98,11 @@ completion for verbs, flags and addresses. Add it to your shell rc: eval "$(libby completion bash)" ``` -Completing an address discovers live daemons the same way `list` does, so -TAB offers daemons after `.` and that daemon's keywords after +Completing an address lists live daemons the same way `list` does, so TAB +offers daemons after `.` and that daemon's keywords after `..`. Results are cached for 10s in `~/.libby/completion_cache.json` so a burst of TABs costs one broadcast, and -discovery is bounded at 0.5s so TAB never hangs. An unreachable broker +the lookup is bounded at 0.5s so TAB never hangs. An unreachable broker completes nothing rather than erroring. ## Modify syntax diff --git a/libby/cli/completion.py b/libby/cli/completion.py index cea5bca..309a7a6 100644 --- a/libby/cli/completion.py +++ b/libby/cli/completion.py @@ -1,4 +1,8 @@ -"""Shell completion for libby addresses, fed by one cached broadcast ``keys.list``.""" +"""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 @@ -17,7 +21,7 @@ @dataclass(frozen=True) class CompletionCache: - """Last discovery result on disk, so a burst of TABs costs one broadcast.""" + """Last peer listing on disk, so a burst of TABs costs one broadcast.""" path: Path = DEFAULT_CACHE_PATH ttl_s: float = CACHE_TTL_S diff --git a/libby/cli/libby_cli.py b/libby/cli/libby_cli.py index 749d3e6..604201e 100644 --- a/libby/cli/libby_cli.py +++ b/libby/cli/libby_cli.py @@ -361,7 +361,7 @@ def cmd_list(namespace: argparse.Namespace) -> int: # Parse before opening a transport so a malformed pattern stays an argument error address = parse_address_pattern(namespace.pattern) # A pattern that spans daemons is answered by broadcast, which always runs - # to its timeout, so it gets the shorter discovery default + # to its timeout, so it gets the shorter broadcast default one_daemon = address.keyword is not None and not address.spans_peers default_timeout = DEFAULT_TIMEOUT_S if one_daemon else DEFAULT_BROADCAST_TIMEOUT_S timeout = namespace.timeout if namespace.timeout is not None else default_timeout @@ -595,7 +595,7 @@ def _connection_key(namespace: argparse.Namespace, config: Dict[str, Any]) -> st return f"zmq {sorted(resolve_address_book(config, namespace.addr).items())}" -def _discover_listings(namespace: argparse.Namespace) -> Listings: +def _fetch_listings(namespace: argparse.Namespace) -> Listings: config = load_cli_config(namespace.config) def fetch(timeout_s: float) -> Listings: @@ -612,7 +612,7 @@ def fetch(timeout_s: float) -> Listings: def _complete_address(prefix: str, parsed_args: argparse.Namespace, **_: Any) -> List[str]: """Complete a partial .. argument.""" try: - return address_candidates(prefix, _discover_listings(parsed_args)) + return address_candidates(prefix, _fetch_listings(parsed_args)) except Exception: # pylint: disable=broad-exception-caught # A completer must never break the shell; an unreachable broker completes nothing return [] @@ -621,7 +621,7 @@ def _complete_address(prefix: str, parsed_args: argparse.Namespace, **_: Any) -> def _complete_peer(prefix: str, parsed_args: argparse.Namespace, **_: Any) -> List[str]: """Complete a partial . argument.""" try: - return peer_candidates(prefix, _discover_listings(parsed_args)) + return peer_candidates(prefix, _fetch_listings(parsed_args)) except Exception: # pylint: disable=broad-exception-caught return [] diff --git a/tests/test_client_integration.py b/tests/test_client_integration.py index ee304d5..b6c14ec 100644 --- a/tests/test_client_integration.py +++ b/tests/test_client_integration.py @@ -25,7 +25,7 @@ GROUP_ID = "hsfei" OTHER_GROUP_ID = "hscal" RPC_TIMEOUT_S = 6.0 -DISCOVERY_TIMEOUT_S = 2.0 +LISTING_TIMEOUT_S = 2.0 SERVE_STOP_TIMEOUT_S = 10.0 INITIAL_POSITION = 10.0 @@ -225,8 +225,8 @@ def test_read_merges_chunked_requests(self): self.assertEqual(list(values), wanted) self.assertTrue(all(entry["ok"] for entry in values.values())) - class DiscoveryCases(unittest.TestCase): - """Broadcast discovery that must hold identically on every transport. + class PeerListingCases(unittest.TestCase): + """Broadcast peer listing that must hold identically on every transport. Concrete subclasses start two fixture daemons in ``GROUP_ID`` and one in ``OTHER_GROUP_ID``, and supply ``client`` plus their qualified ids. @@ -238,27 +238,27 @@ class DiscoveryCases(unittest.TestCase): def test_peers_lists_every_daemon_in_the_group(self): """Find the group's daemons by wildcard, and no other group's.""" - found = self.client.peers(f"{GROUP_ID}.%", timeout_s=DISCOVERY_TIMEOUT_S) + found = self.client.peers(f"{GROUP_ID}.%", timeout_s=LISTING_TIMEOUT_S) for peer in self.group_peers: self.assertIn(peer, found) self.assertNotIn(self.other_peer, found) def test_peers_spans_groups_and_skips_the_client(self): """Find every daemon with %.% without listing the asking client.""" - found = self.client.peers("%.%", timeout_s=DISCOVERY_TIMEOUT_S) + found = self.client.peers("%.%", timeout_s=LISTING_TIMEOUT_S) for peer in (*self.group_peers, self.other_peer): self.assertIn(peer, found) self.assertNotIn(DEFAULT_SELF_ID, found) def test_peers_with_an_exact_id_confirms_one_daemon(self): """Resolve an exact . to just that daemon.""" - found = self.client.peers(self.group_peers[0], timeout_s=DISCOVERY_TIMEOUT_S) + found = self.client.peers(self.group_peers[0], timeout_s=LISTING_TIMEOUT_S) self.assertEqual(found, [self.group_peers[0]]) def test_list_across_daemons_returns_qualified_names(self): """List one keyword on every daemon of a group as names that feed back in.""" names = self.client.list(f"{GROUP_ID}.%.positionvalue", - timeout_s=DISCOVERY_TIMEOUT_S) + timeout_s=LISTING_TIMEOUT_S) for peer in self.group_peers: self.assertIn(f"{peer}.positionvalue", names) self.assertNotIn(f"{self.other_peer}.positionvalue", names) @@ -270,12 +270,12 @@ def test_list_across_daemons_returns_qualified_names(self): def test_list_across_daemons_with_no_match_is_empty(self): """Return nothing, rather than raise, when no daemon has the keyword.""" self.assertEqual( - self.client.list(f"{GROUP_ID}.%.nosuchkeyword", timeout_s=DISCOVERY_TIMEOUT_S), + self.client.list(f"{GROUP_ID}.%.nosuchkeyword", timeout_s=LISTING_TIMEOUT_S), []) def test_peer_listings_carry_each_daemons_keywords(self): """Map each daemon to its keyword names from the same broadcast.""" - listings = self.client.peer_listings(f"{GROUP_ID}.%", timeout_s=DISCOVERY_TIMEOUT_S) + listings = self.client.peer_listings(f"{GROUP_ID}.%", timeout_s=LISTING_TIMEOUT_S) for peer in self.group_peers: self.assertIn("positionvalue", listings[peer]) self.assertIn("uptime", listings[peer]) @@ -358,20 +358,20 @@ def tearDownClass(cls): @unittest.skipUnless(_broker_available(), "no RabbitMQ broker reachable at amqp://localhost") -class RabbitMQDiscoveryTests(_Bases.DiscoveryCases): - """Discovery cases over RabbitMQ.""" +class RabbitMQPeerListingTests(_Bases.PeerListingCases): + """Peer listing cases over RabbitMQ.""" daemons: List[LibbyDaemon] @classmethod def setUpClass(cls): cls.daemons = [ - _started_daemon(_RabbitFixtureDaemon, "discoveryonermq", GROUP_ID), - _started_daemon(_RabbitFixtureDaemon, "discoverytwormq", GROUP_ID), - _started_daemon(_RabbitFixtureDaemon, "discoveryotherrmq", OTHER_GROUP_ID), + _started_daemon(_RabbitFixtureDaemon, "listingonermq", GROUP_ID), + _started_daemon(_RabbitFixtureDaemon, "listingtwormq", GROUP_ID), + _started_daemon(_RabbitFixtureDaemon, "listingotherrmq", OTHER_GROUP_ID), ] - cls.group_peers = (f"{GROUP_ID}.discoveryonermq", f"{GROUP_ID}.discoverytwormq") - cls.other_peer = f"{OTHER_GROUP_ID}.discoveryotherrmq" + cls.group_peers = (f"{GROUP_ID}.listingonermq", f"{GROUP_ID}.listingtwormq") + cls.other_peer = f"{OTHER_GROUP_ID}.listingotherrmq" cls.client = Client.rabbitmq(rabbitmq_url=RABBITMQ_URL) @classmethod @@ -381,22 +381,22 @@ def tearDownClass(cls): daemon.stop() -class ZmqDiscoveryTests(_Bases.DiscoveryCases): - """Discovery cases over ZMQ, where the address book is what gets asked.""" +class ZmqPeerListingTests(_Bases.PeerListingCases): + """Peer listing cases over ZMQ, where the address book is what gets asked.""" daemons: List[LibbyDaemon] @classmethod def setUpClass(cls): - groups = {"discoveryonezmq": GROUP_ID, "discoverytwozmq": GROUP_ID, - "discoveryotherzmq": OTHER_GROUP_ID} + groups = {"listingonezmq": GROUP_ID, "listingtwozmq": GROUP_ID, + "listingotherzmq": OTHER_GROUP_ID} address_book = {f"{group}.{name}": _free_endpoint() for name, group in groups.items()} cls.daemons = [ _started_daemon(_ZmqFixtureDaemon, name, group, bind=address_book[f"{group}.{name}"]) for name, group in groups.items() ] - cls.group_peers = (f"{GROUP_ID}.discoveryonezmq", f"{GROUP_ID}.discoverytwozmq") - cls.other_peer = f"{OTHER_GROUP_ID}.discoveryotherzmq" + cls.group_peers = (f"{GROUP_ID}.listingonezmq", f"{GROUP_ID}.listingtwozmq") + cls.other_peer = f"{OTHER_GROUP_ID}.listingotherzmq" cls.client = Client.zmq(bind=_free_endpoint(), address_book=address_book) @classmethod From bbb14aba33d6480cc5ff180f79a1b991ba4e6885 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:49:04 -0700 Subject: [PATCH 4/6] State that peer discovery is not implemented rather than describing its partial behavior --- README.md | 22 ++++++++++++---------- 1 file changed, 12 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index 72c7ba1..eb32903 100644 --- a/README.md +++ b/README.md @@ -39,18 +39,20 @@ 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. -Bamboo's own hello/discovery is a separate mechanism and is not a way to find -peers today: - -- `Protocol` never instantiates a `PeerTable`, so `Libby.peers_alive()` returns - `{}` and `Libby.wait_for_peer()` always fails, on either transport. -- `Libby.rabbitmq()` doesn't start discovery at all. The broker routes messages; - it does not tell a peer who else is connected. -- On ZMQ a hello only reaches peers already in the sender's address book, so a - client learns nothing about a daemon that hasn't been told about the client. +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 that mechanism, not +`LibbyDaemon`'s `discovery_enabled` / `on_hello` control bamboo's hello, not the broadcast above. ## Testing From 6ce1003853cfd2834ab6fe46c0df9c6644b1f808 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Tue, 22 Sep 2026 13:24:17 -0700 Subject: [PATCH 5/6] List only daemons as peers, since every libby connection answers keys.list --- docs/source/cli.md | 4 ++ docs/source/client.md | 6 +++ libby/client.py | 14 +++++-- libby/daemon.py | 2 + libby/libby.py | 10 +++++ tests/test_client_integration.py | 19 ++++++++- tests/test_peer_listings.py | 72 ++++++++++++++++++++++++++++++++ 7 files changed, 123 insertions(+), 4 deletions(-) create mode 100644 tests/test_peer_listings.py diff --git a/docs/source/cli.md b/docs/source/cli.md index 79750b6..fd9b653 100644 --- a/docs/source/cli.md +++ b/docs/source/cli.md @@ -86,6 +86,10 @@ 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. diff --git a/docs/source/client.md b/docs/source/client.md index b1f0b44..047f20d 100644 --- a/docs/source/client.md +++ b/docs/source/client.md @@ -134,5 +134,11 @@ 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 the rest are filtered out. +A peer running a libby from before that flag omits it and is still listed, +so this does not hide daemons that have yet to be redeployed. + See {mod}`libby.client` in the {doc}`API reference ` for the full method signatures. diff --git a/libby/client.py b/libby/client.py index a85910e..6c98415 100644 --- a/libby/client.py +++ b/libby/client.py @@ -200,16 +200,24 @@ def peer_listings( A keyword segment in ``pattern`` narrows the names; without one every keyword is listed. Same timeout semantics as :meth:`peers`. + + Clients answer the broadcast too, since every ``Libby`` serves + ``keys.list``, so a reply counts only if it does not report itself as + a non-daemon. A peer on a libby predating that flag omits it and is + still listed; only a peer on this libby can broadcast at all, and + those always report it. """ address = parse_address_pattern(pattern) replies = self._libby.broadcast_request( "keys.list", {"pattern": address.keyword or "%"}, timeout_s=timeout_s) + daemons = [reply for reply in replies + if reply.payload.get("ok") and reply.payload.get("is_daemon", True)] wanted = set(match_pattern(address.peer_pattern, - (reply.peer_id for reply in replies))) + (reply.peer_id for reply in daemons))) return { reply.peer_id: list(reply.payload.get("matches", [])) - for reply in replies - if reply.peer_id in wanted and reply.payload.get("ok") + for reply in daemons + if reply.peer_id in wanted } def describe(self, name: str, *, timeout_s: float = DEFAULT_TIMEOUT_S) -> Dict[str, Any]: diff --git a/libby/daemon.py b/libby/daemon.py index 75482fb..4ab9e2b 100644 --- a/libby/daemon.py +++ b/libby/daemon.py @@ -361,6 +361,7 @@ def build_libby(self) -> Libby: keys=[], callback=None, group_id=self.config_group_id(), + is_daemon=True, ) if transport == "zmq": @@ -374,6 +375,7 @@ def build_libby(self) -> Libby: discover_interval_s=self.config_discovery_interval_s(), hello_on_start=True, group_id=self.config_group_id(), + is_daemon=True, ) raise ValueError( diff --git a/libby/libby.py b/libby/libby.py index 695480d..5f6664a 100644 --- a/libby/libby.py +++ b/libby/libby.py @@ -36,9 +36,13 @@ def __init__( discover: bool = False, discover_interval_s: float = 5.0, hello_on_start: bool = True, + is_daemon: bool = False, ): self.self_id = self_id self.transport = transport + # Every Libby answers keys.list, so a broadcast reaches clients too. + # Only LibbyDaemon sets this, and only these are listed as peers + self.is_daemon = is_daemon self.keys = KeyRegistry() self.proto = Protocol(transport=self.transport, self_id=self_id, keys=self.keys) @@ -80,6 +84,7 @@ def zmq( discover_interval_s: float = 5.0, hello_on_start: bool = True, group_id: Optional[str] = None, + is_daemon: bool = False, ) -> "Libby": try: from .zmq_transport import ZmqTransport @@ -101,6 +106,7 @@ def zmq( discover=discover, discover_interval_s=discover_interval_s, hello_on_start=hello_on_start, + is_daemon=is_daemon, ) @classmethod @@ -111,6 +117,8 @@ def rabbitmq( keys: Optional[List[str]] = None, callback: Optional[Callable[[dict, dict], Optional[dict]]] = None, group_id: Optional[str] = None, + *, + is_daemon: bool = False, ) -> "Libby": """ Create a Libby instance using RabbitMQ transport. @@ -154,6 +162,7 @@ def rabbitmq( discover=False, discover_interval_s=0, hello_on_start=False, + is_daemon=is_daemon, ) # lifecycle @@ -260,6 +269,7 @@ def _keys_list(self, payload: dict, _ctx: dict) -> dict: "ok": True, "matches": match_pattern(pattern, self._keywords), "services": self._served_services(), + "is_daemon": self.is_daemon, } def _served_services(self) -> List[str]: diff --git a/tests/test_client_integration.py b/tests/test_client_integration.py index b6c14ec..a9e47a2 100644 --- a/tests/test_client_integration.py +++ b/tests/test_client_integration.py @@ -229,12 +229,15 @@ class PeerListingCases(unittest.TestCase): """Broadcast peer listing that must hold identically on every transport. Concrete subclasses start two fixture daemons in ``GROUP_ID`` and one - in ``OTHER_GROUP_ID``, and supply ``client`` plus their qualified ids. + in ``OTHER_GROUP_ID``, plus a plain ``Client`` whose ``self_id`` is + shaped like a peer in ``GROUP_ID``, and supply ``client`` plus the + qualified ids. """ client: Client group_peers: Tuple[str, str] other_peer: str + impostor_peer: str def test_peers_lists_every_daemon_in_the_group(self): """Find the group's daemons by wildcard, and no other group's.""" @@ -250,6 +253,12 @@ def test_peers_spans_groups_and_skips_the_client(self): self.assertIn(peer, found) self.assertNotIn(DEFAULT_SELF_ID, found) + def test_peers_omits_a_client_that_answers_the_broadcast(self): + """Leave out a plain Client, which serves keys.list like a daemon does.""" + found = self.client.peers(f"{GROUP_ID}.%", timeout_s=LISTING_TIMEOUT_S) + self.assertNotIn(self.impostor_peer, found) + self.assertIn(self.group_peers[0], found) + def test_peers_with_an_exact_id_confirms_one_daemon(self): """Resolve an exact . to just that daemon.""" found = self.client.peers(self.group_peers[0], timeout_s=LISTING_TIMEOUT_S) @@ -372,11 +381,14 @@ def setUpClass(cls): ] cls.group_peers = (f"{GROUP_ID}.listingonermq", f"{GROUP_ID}.listingtwormq") cls.other_peer = f"{OTHER_GROUP_ID}.listingotherrmq" + cls.impostor_peer = f"{GROUP_ID}.impostorrmq" + cls.impostor = Client.rabbitmq(self_id=cls.impostor_peer, rabbitmq_url=RABBITMQ_URL) cls.client = Client.rabbitmq(rabbitmq_url=RABBITMQ_URL) @classmethod def tearDownClass(cls): cls.client.close() + cls.impostor.close() for daemon in cls.daemons: daemon.stop() @@ -397,11 +409,16 @@ def setUpClass(cls): ] cls.group_peers = (f"{GROUP_ID}.listingonezmq", f"{GROUP_ID}.listingtwozmq") cls.other_peer = f"{OTHER_GROUP_ID}.listingotherzmq" + cls.impostor_peer = f"{GROUP_ID}.impostorzmq" + impostor_endpoint = _free_endpoint() + address_book[cls.impostor_peer] = impostor_endpoint + cls.impostor = Client.zmq(self_id=cls.impostor_peer, bind=impostor_endpoint) cls.client = Client.zmq(bind=_free_endpoint(), address_book=address_book) @classmethod def tearDownClass(cls): cls.client.close() + cls.impostor.close() for daemon in cls.daemons: daemon.stop() diff --git a/tests/test_peer_listings.py b/tests/test_peer_listings.py new file mode 100644 index 0000000..dc525fc --- /dev/null +++ b/tests/test_peer_listings.py @@ -0,0 +1,72 @@ +"""Unit tests for which broadcast replies ``Client.peer_listings`` keeps. + +Every ``Libby`` serves ``keys.list``, so clients answer the broadcast too and +have to be filtered out. The reply shapes are produced by a stand-in here; the +end-to-end path over both transports lives in test_client_integration. +""" +from __future__ import annotations + +import unittest +from typing import Any, Dict, List + +from libby import Client +from libby.libby import BroadcastReply + + +def _reply(peer_id: str, **payload: Any) -> BroadcastReply: + return BroadcastReply(peer_id, {"ok": True, "matches": ["uptime"], **payload}) + + +class _BroadcastLibby: # pylint: disable=too-few-public-methods + """Stands in for Libby, answering a broadcast with canned replies.""" + + def __init__(self, replies: List[BroadcastReply]): + self._replies = replies + + def broadcast_request(self, key: str, payload: Dict[str, Any], + timeout_s: float = 1.0) -> List[BroadcastReply]: + """Answer with the canned replies.""" + # Arguments are ignored; the signature exists to match Libby + # pylint: disable=unused-argument + return self._replies + + +class PeerListingFilterTests(unittest.TestCase): + """Which replies count as a daemon.""" + + def _peers(self, *replies: BroadcastReply) -> List[str]: + return Client(_BroadcastLibby(list(replies))).peers("hsfei.%") + + def test_keeps_a_daemon(self): + self.assertEqual(self._peers(_reply("hsfei.adc", is_daemon=True)), ["hsfei.adc"]) + + def test_drops_a_client(self): + self.assertEqual(self._peers(_reply("hsfei.impostor", is_daemon=False)), []) + + def test_keeps_a_reply_predating_the_flag(self): + # An older libby omits is_daemon; only this libby can broadcast at all, + # so a missing flag means a daemon that has not been redeployed yet + self.assertEqual(self._peers(_reply("hsfei.adc")), ["hsfei.adc"]) + + def test_drops_a_failed_reply(self): + self.assertEqual( + self._peers(BroadcastReply("hsfei.adc", {"ok": False, "error": "nope"})), []) + + def test_filters_before_matching_the_pattern(self): + found = self._peers( + _reply("hsfei.adc", is_daemon=True), + _reply("hsfei.impostor", is_daemon=False), + _reply("hscal.hkettherm", is_daemon=True), + ) + self.assertEqual(found, ["hsfei.adc"]) + + def test_peer_listings_carry_the_matches(self): + listings = Client(_BroadcastLibby([ + _reply("hsfei.adc", is_daemon=True, matches=["uptime", "isconnected"]), + _reply("hsfei.impostor", is_daemon=False, matches=["uptime"]), + ])).peer_listings("hsfei.%") + self.assertEqual(listings, {"hsfei.adc": ["uptime", "isconnected"]}) + + +if __name__ == "__main__": + unittest.main() From 69db4d5f76615e5f2e968ab475bc2c040b2a0935 Mon Sep 17 00:00:00 2001 From: Mike Langmayr <1809691+mikelangmayr@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:24:09 -0700 Subject: [PATCH 6/6] Require an explicit daemon flag before listing a peer --- docs/source/client.md | 6 +++--- libby/client.py | 8 +++----- tests/test_peer_listings.py | 7 +++---- 3 files changed, 9 insertions(+), 12 deletions(-) diff --git a/docs/source/client.md b/docs/source/client.md index 047f20d..0159ae4 100644 --- a/docs/source/client.md +++ b/docs/source/client.md @@ -136,9 +136,9 @@ 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 the rest are filtered out. -A peer running a libby from before that flag omits it and is still listed, -so this does not hide daemons that have yet to be redeployed. +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 ` for the full method signatures. diff --git a/libby/client.py b/libby/client.py index 6c98415..0d6b1b0 100644 --- a/libby/client.py +++ b/libby/client.py @@ -202,16 +202,14 @@ def peer_listings( keyword is listed. Same timeout semantics as :meth:`peers`. Clients answer the broadcast too, since every ``Libby`` serves - ``keys.list``, so a reply counts only if it does not report itself as - a non-daemon. A peer on a libby predating that flag omits it and is - still listed; only a peer on this libby can broadcast at all, and - those always report it. + ``keys.list``, so a reply counts only if it reports itself as a + daemon. """ address = parse_address_pattern(pattern) replies = self._libby.broadcast_request( "keys.list", {"pattern": address.keyword or "%"}, timeout_s=timeout_s) daemons = [reply for reply in replies - if reply.payload.get("ok") and reply.payload.get("is_daemon", True)] + if reply.payload.get("ok") and reply.payload.get("is_daemon")] wanted = set(match_pattern(address.peer_pattern, (reply.peer_id for reply in daemons))) return { diff --git a/tests/test_peer_listings.py b/tests/test_peer_listings.py index dc525fc..58e2baf 100644 --- a/tests/test_peer_listings.py +++ b/tests/test_peer_listings.py @@ -43,10 +43,9 @@ def test_keeps_a_daemon(self): def test_drops_a_client(self): self.assertEqual(self._peers(_reply("hsfei.impostor", is_daemon=False)), []) - def test_keeps_a_reply_predating_the_flag(self): - # An older libby omits is_daemon; only this libby can broadcast at all, - # so a missing flag means a daemon that has not been redeployed yet - self.assertEqual(self._peers(_reply("hsfei.adc")), ["hsfei.adc"]) + def test_drops_a_reply_predating_the_flag(self): + # A libby old enough to omit is_daemon cannot prove it is one + self.assertEqual(self._peers(_reply("hsfei.adc")), []) def test_drops_a_failed_reply(self): self.assertEqual(