From c540051f94f7e90c08222ecb9c57b3a1e1db7e0b Mon Sep 17 00:00:00 2001 From: Prakriti Gupta Date: Wed, 22 Jul 2026 16:49:01 -0700 Subject: [PATCH] Update LibbyDaemon base class --- libby/config.py | 280 ++++++++++++++++++----- libby/daemon.py | 581 +++++++++++++++++++++++++++++++++++------------- peers/peer_a.py | 47 ++-- peers/peer_b.py | 44 ++-- peers/peer_c.py | 50 ++--- 5 files changed, 723 insertions(+), 279 deletions(-) diff --git a/libby/config.py b/libby/config.py index ea89a32..bc8e57b 100644 --- a/libby/config.py +++ b/libby/config.py @@ -1,64 +1,238 @@ from __future__ import annotations -from typing import Any, Dict, Mapping -import json, os, pathlib + +import copy +import json +import os +import pathlib +from typing import Any, Dict, List, Mapping, Optional, Union try: import yaml except Exception: yaml = None -def _load_json(p: pathlib.Path) -> Dict[str, Any]: - return json.loads(p.read_text()) -def _load_yaml(p: pathlib.Path) -> Dict[str, Any]: +PathLike = Union[str, os.PathLike] + + +class ConfigError(ValueError): + """Raised when daemon configuration cannot be selected or validated.""" + + +def _load_json(path: pathlib.Path) -> Dict[str, Any]: + parsed = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(parsed, dict): + raise ConfigError(f"configuration root must be a mapping: {path}") + return parsed + + +def _load_yaml(path: pathlib.Path) -> Dict[str, Any]: if yaml is None: - raise RuntimeError("YAML requested but PyYAML not installed. `pip install pyyaml`") - return yaml.safe_load(p.read_text()) or {} - -def load_config(path: str | os.PathLike[str]) -> Dict[str, Any]: - """ - Load config from .json or .yml/.yaml. - If the extension is missing/unknown, attempt JSON → YAML. - """ - p = pathlib.Path(path) - if not p.exists(): - raise FileNotFoundError(f"Config file not found: {p}") - ext = p.suffix.lower() - if ext == ".json": - return _load_json(p) or {} - if ext in (".yml", ".yaml"): - return _load_yaml(p) or {} - # Auto-detect - for fn in (_load_json, _load_yaml): + raise RuntimeError( + "YAML requested but PyYAML is not installed; " + "install it with `pip install pyyaml`" + ) + + parsed = yaml.safe_load(path.read_text(encoding="utf-8")) or {} + if not isinstance(parsed, dict): + raise ConfigError(f"configuration root must be a mapping: {path}") + return parsed + + +def load_config(path: PathLike) -> Dict[str, Any]: + """Load a JSON or YAML configuration file.""" + + config_path = pathlib.Path(path) + if not config_path.exists(): + raise FileNotFoundError(f"config file not found: {config_path}") + + extension = config_path.suffix.lower() + if extension == ".json": + return _load_json(config_path) + if extension in (".yml", ".yaml"): + return _load_yaml(config_path) + + errors: List[Exception] = [] + for loader in (_load_json, _load_yaml): try: - return fn(p) or {} - except Exception: - pass - raise ValueError(f"Could not parse config file as JSON or YAML: {p}") - -def with_env_overrides(cfg: Mapping[str, Any], prefix: str = "LIBBY_") -> Dict[str, Any]: - """ - Uppercase, underscore keys: LIBBY_PEER_ID, LIBBY_BIND, etc. - Booleans: '1','true','yes' => True ; '0','false','no' => False - Lists: comma-separated. - """ - out: Dict[str, Any] = dict(cfg) - - def coerce(v: str) -> Any: - s = v.strip() - ls = s.lower() - if ls in ("true","1","yes","on"): return True - if ls in ("false","0","no","off"): return False - if "," in s: return [x.strip() for x in s.split(",")] + return loader(config_path) + except Exception as exc: + errors.append(exc) + + raise ConfigError( + f"could not parse config file as JSON or YAML: {config_path}" + ) from errors[-1] + + +def with_env_overrides( + config: Mapping[str, Any], + prefix: str = "LIBBY_", +) -> Dict[str, Any]: + """Apply top-level environment-variable overrides.""" + + result = dict(config) + + def coerce(value: str) -> Any: + stripped = value.strip() + lowered = stripped.lower() + + if lowered in ("true", "1", "yes", "on"): + return True + if lowered in ("false", "0", "no", "off"): + return False + if "," in stripped: + return [item.strip() for item in stripped.split(",")] + try: - if "." in s: return float(s) - return int(s) - except Exception: - return s - - for k, v in os.environ.items(): - if not k.startswith(prefix): - continue - key = k[len(prefix):].lower() - out[key] = coerce(v) - return out + if "." in stripped: + return float(stripped) + return int(stripped) + except ValueError: + return stripped + + for env_name, env_value in os.environ.items(): + if env_name.startswith(prefix): + key = env_name[len(prefix):].lower() + result[key] = coerce(env_value) + + return result + + +def deep_merge( + base: Mapping[str, Any], + override: Mapping[str, Any], +) -> Dict[str, Any]: + """Return a recursive merge without mutating either input.""" + + result: Dict[str, Any] = copy.deepcopy(dict(base)) + + for key, value in override.items(): + existing = result.get(key) + if isinstance(existing, Mapping) and isinstance(value, Mapping): + result[key] = deep_merge(existing, value) + else: + result[key] = copy.deepcopy(value) + + return result + + +def is_subsystem_config( + config: Mapping[str, Any], + *, + daemon_section: str = "daemons", +) -> bool: + """Return whether a config contains a daemon collection.""" + + return daemon_section in config + + +def list_daemons( + config: Mapping[str, Any], + *, + daemon_section: str = "daemons", +) -> List[str]: + """List daemon IDs in a subsystem config.""" + + daemons = config.get(daemon_section, {}) + if not isinstance(daemons, Mapping): + raise ConfigError( + f"{daemon_section!r} must contain a mapping of daemon IDs" + ) + return list(daemons) + + +def extract_daemon_config( + full_config: Mapping[str, Any], + daemon_id: str, + *, + daemon_section: str = "daemons", +) -> Dict[str, Any]: + """Merge subsystem defaults with one daemon's overrides.""" + + daemons = full_config.get(daemon_section, {}) + if not isinstance(daemons, Mapping): + raise ConfigError( + f"{daemon_section!r} must contain a mapping of daemon IDs" + ) + + if daemon_id not in daemons: + raise ConfigError( + f"daemon {daemon_id!r} is not defined; " + f"available daemons: {list(daemons)}" + ) + + daemon_override = daemons[daemon_id] or {} + if not isinstance(daemon_override, Mapping): + raise ConfigError( + f"configuration for daemon {daemon_id!r} must be a mapping" + ) + + defaults = { + key: value + for key, value in full_config.items() + if key != daemon_section + } + result = deep_merge(defaults, daemon_override) + result.setdefault("peer_id", daemon_id) + return result + + +class DaemonConfigLoader: + """Load a single-daemon or multi-daemon subsystem config.""" + + def __init__( + self, + path: PathLike, + *, + daemon_section: str = "daemons", + ) -> None: + self.path = pathlib.Path(path) + self.daemon_section = daemon_section + self._config: Optional[Dict[str, Any]] = None + + @property + def config(self) -> Dict[str, Any]: + if self._config is None: + self._config = load_config(self.path) + return self._config + + @property + def is_subsystem(self) -> bool: + return is_subsystem_config( + self.config, + daemon_section=self.daemon_section, + ) + + @property + def subsystem(self) -> Optional[str]: + raw = self.config.get("subsystem") + return str(raw) if raw is not None else None + + @property + def daemon_ids(self) -> List[str]: + if self.is_subsystem: + return list_daemons( + self.config, + daemon_section=self.daemon_section, + ) + + peer_id = self.config.get("peer_id", self.path.stem) + return [str(peer_id)] + + def get_daemon_config( + self, + daemon_id: Optional[str] = None, + ) -> Dict[str, Any]: + if self.is_subsystem: + if daemon_id is None: + raise ConfigError( + "daemon_id is required for a subsystem config; " + f"available daemons: {self.daemon_ids}" + ) + return extract_daemon_config( + self.config, + daemon_id, + daemon_section=self.daemon_section, + ) + + return copy.deepcopy(self.config) \ No newline at end of file diff --git a/libby/daemon.py b/libby/daemon.py index 7d7d5da..02cf18b 100644 --- a/libby/daemon.py +++ b/libby/daemon.py @@ -1,162 +1,322 @@ from __future__ import annotations -import json -from dataclasses import is_dataclass, asdict + import collections.abc as cabc -import signal, sys, threading, time -from typing import Any, Callable, Dict, List, Optional +import json +import logging +import signal +import threading +import time +from dataclasses import asdict, is_dataclass +from typing import Any, Callable, Dict, Iterable, List, Mapping, Optional, Type, TypeVar + +from .config import ConfigError, DaemonConfigLoader, with_env_overrides +from .keyword import Keyword from .libby import Libby + Payload = Dict[str, Any] -RPCHandler = Callable[[Payload], Dict[str, Any]] +RPCHandler = Callable[[Payload], Any] EvtHandler = Callable[[Payload], None] +DaemonT = TypeVar("DaemonT", bound="LibbyDaemon") + class LibbyDaemon: - """ - Base daemon class for Libby peers with support for multiple transports. + """Base class for configurable Libby daemons. - ZMQ Usage: - class MyPeer(LibbyDaemon): - peer_id = "my-peer" - bind = "tcp://*:5555" - address_book = {"other-peer": "tcp://localhost:5556"} + The class owns the following: - services = {"echo": lambda payload: {"echo": payload}} - topics = {"alerts": lambda payload: print(payload)} + * Libby transport construction + * service and topic registration + * typed keyword registration + * YAML/JSON configuration loading + * subsystem configuration selection + * logging and lifecycle management - RabbitMQ Usage: - class MyPeer(LibbyDaemon): - transport = "rabbitmq" - peer_id = "my-peer" - rabbitmq_url = "amqp://localhost" # optional, defaults to this + Instrument-specific subclasses should normally only implement ``on_start`` + and ``on_stop``, plus any hardware-facing methods. + """ - services = {"echo": lambda payload: {"echo": payload}} - topics = {"alerts": lambda payload: print(payload)} + CONFIG_ATTRIBUTES = frozenset( + { + "peer_id", + "bind", + "address_book", + "discovery_enabled", + "discovery_interval_s", + "transport", + "rabbitmq_url", + "group_id", + "fail_fast_on_start", + } + ) - Note: RabbitMQ doesn't need bind or address_book since routing is - handled automatically by the broker. - """ - # simple attributes users set peer_id: Optional[str] = None bind: Optional[str] = None address_book: Optional[Dict[str, str]] = None discovery_enabled: bool = True discovery_interval_s: float = 5.0 - # transport selection: "zmq" (default) or "rabbitmq" + # Keep Libby's existing default. Instrument configurations can select + # their preferred transport explicitly without changing behavior for existing ZMQ users. transport: str = "zmq" rabbitmq_url: Optional[str] = None group_id: Optional[str] = None - # internal config - _config: Dict[str, Any] = {} - # payload-only handlers + + # A daemon that fails to initialize should not continue advertising a + # partially initialized service. + fail_fast_on_start: bool = True + services: Dict[str, RPCHandler] = {} topics: Dict[str, EvtHandler] = {} def __init__(self) -> None: - self.services: Dict[str, RPCHandler] = dict(getattr(type(self), "services", {})) - self.topics: Dict[str, EvtHandler] = dict(getattr(type(self), "topics", {})) + # Copy subclass declarations so instances never mutate class-level maps. + self.services = dict(getattr(type(self), "services", {})) + self.topics = dict(getattr(type(self), "topics", {})) + + self._config: Dict[str, Any] = {} + self._pending_keywords: List[Keyword] = [] + self._stop_event = threading.Event() + self._started = False + self.libby: Optional[Libby] = None + self.logger = logging.getLogger(type(self).__name__) @classmethod - def from_config(cls, cfg: Dict[str, Any]) -> "LibbyDaemon": - d = cls() - def set_if(k: str): - v = cfg.get(k) - if v not in (None, "", {}, []): - setattr(d, k, v) - - set_if("peer_id") - set_if("bind") - set_if("address_book") - set_if("discovery_enabled") - set_if("discovery_interval_s") - - if isinstance(cfg.get("services"), dict): - d.services.update(cfg["services"]) - if isinstance(cfg.get("topics"), dict): - d.topics.update(cfg["topics"]) - return d + def config_attributes(cls) -> frozenset[str]: + """Return config keys mapped directly onto daemon attributes. + + Subclasses may extend the set: + + CONFIG_ATTRIBUTES = ( + LibbyDaemon.CONFIG_ATTRIBUTES | {"device_name"} + ) + """ + + return cls.CONFIG_ATTRIBUTES @classmethod - def from_config_file(cls, path: str) -> "LibbyDaemon": - import os, json - try: - import yaml # type: ignore - except Exception: - yaml = None + def from_config( + cls: Type[DaemonT], + config: Mapping[str, Any], + ) -> DaemonT: + """Build a daemon from a configuration mapping.""" + + if not isinstance(config, Mapping): + raise TypeError("daemon configuration must be a mapping") - with open(path, "r", encoding="utf-8") as f: - text = f.read() + instance = cls() + instance._config = dict(config) - if path.endswith((".yml", ".yaml")) and yaml is not None: - cfg = yaml.safe_load(text) or {} + for attribute in instance.config_attributes(): + if attribute in config: + setattr(instance, attribute, config[attribute]) + + configured_services = config.get("services") + if isinstance(configured_services, Mapping): + instance.services.update(configured_services) + + configured_topics = config.get("topics") + if isinstance(configured_topics, Mapping): + instance.topics.update(configured_topics) + + instance._setup_logging() + return instance + + @classmethod + def from_config_file( + cls: Type[DaemonT], + path: str, + daemon_id: Optional[str] = None, + *, + env_prefix: Optional[str] = None, + ) -> DaemonT: + """Build a daemon from a single-daemon or subsystem config file. + + For a subsystem config containing one daemon, the daemon is selected + automatically. For multiple daemons, ``daemon_id`` is required. + """ + + loader = DaemonConfigLoader(path) + + if loader.is_subsystem and daemon_id is None: + if len(loader.daemon_ids) == 1: + daemon_id = loader.daemon_ids[0] + else: + raise ConfigError( + "Subsystem config contains multiple daemons: " + f"{loader.daemon_ids}. Specify daemon_id." + ) + + config = loader.get_daemon_config(daemon_id) + if env_prefix: + config = with_env_overrides(config, prefix=env_prefix) + + return cls.from_config(config) + + def get_config(self, key: str, default: Any = None) -> Any: + """Read a configuration value using optional dot notation.""" + + value: Any = self._config + for part in key.split("."): + if not isinstance(value, Mapping) or part not in value: + return default + value = value[part] + return value + + def _setup_logging(self) -> None: + """Configure a daemon-local logger without resetting global logging.""" + + self.logger = logging.getLogger( + self.peer_id or type(self).__name__ + ) + + raw = self._config.get("logging") + if not isinstance(raw, Mapping): + return + + level_name = str(raw.get("level", "INFO")).upper() + level = getattr(logging, level_name, logging.INFO) + self.logger.setLevel(level) + + if "propagate" in raw: + self.logger.propagate = bool(raw["propagate"]) + + # Only install a daemon-owned handler when explicitly requested by a + # logging section. This avoids logging.basicConfig() changing the host + # application's global logging policy + if any(getattr(handler, "_libby_daemon_handler", False) + for handler in self.logger.handlers): + return + + log_file = raw.get("file") + handler: logging.Handler + if log_file: + handler = logging.FileHandler(str(log_file)) else: - cfg = json.loads(text or "{}") - - if not isinstance(cfg, dict): - raise ValueError(f"Config file {path} did not parse to a dict.") - - return cls.from_config(cfg) - - # optional hooks - def on_start(self, libby: Libby) -> None: ... - def on_stop(self, libby: Optional[Libby] = None) -> None: ... - def on_hello(self, libby: Libby) -> None: ... - def on_event(self, topic: str, msg) -> None: - print(f"[{self.__class__.__name__}] {topic}: {msg.env.payload}") - - def config_peer_id(self) -> str: return self.peer_id or self._must("peer_id") - def config_bind(self) -> str: return self.bind or self._must("bind") - def config_address_book(self) -> Dict[str, str]: return self.address_book or self._must("address_book") - def config_rabbitmq_url(self) -> str: return self.rabbitmq_url or "amqp://localhost" - def config_group_id(self) -> Optional[str]: return self.group_id - def config_address_book(self) -> Dict[str, str]: return self.address_book if self.address_book is not None else {} - def config_discovery_enabled(self) -> bool: return bool(self.discovery_enabled) - def config_discovery_interval_s(self) -> float: return float(self.discovery_interval_s) - def config_rpc_keys(self) -> List[str]: return list(self.services.keys()) - def config_subscriptions(self) -> List[str]: return list(self.topics.keys()) - - # user-facing helpers + handler = logging.StreamHandler() + + handler.setLevel(level) + handler.setFormatter( + logging.Formatter( + str( + raw.get( + "format", + "%(asctime)s - %(name)s - " + "%(levelname)s - %(message)s", + ) + ) + ) + ) + setattr(handler, "_libby_daemon_handler", True) + self.logger.addHandler(handler) + + + ### Optional hooks + + def on_start(self, libby: Libby) -> None: + """Initialize hardware and register keywords.""" + + def on_stop(self, libby: Optional[Libby] = None) -> None: + """Release hardware resources.""" + + def on_hello(self, libby: Libby) -> None: + """Run after the initial discovery hello.""" + + def on_event(self, topic: str, msg: Any) -> None: + self.logger.info("%s: %s", topic, msg.env.payload) + + + ### Config accessors + + def config_peer_id(self) -> str: + return self.peer_id or self._must("peer_id") + + def config_bind(self) -> str: + return self.bind or self._must("bind") + + def config_address_book(self) -> Dict[str, str]: + return dict(self.address_book or {}) + + def config_rabbitmq_url(self) -> str: + return self.rabbitmq_url or "amqp://localhost" + + def config_group_id(self) -> Optional[str]: + return self.group_id + + def config_discovery_enabled(self) -> bool: + return bool(self.discovery_enabled) + + def config_discovery_interval_s(self) -> float: + return float(self.discovery_interval_s) + + def config_rpc_keys(self) -> List[str]: + return list(self.services) + + def config_subscriptions(self) -> List[str]: + return list(self.topics) + def add_service(self, key: str, fn: RPCHandler) -> None: self.services[key] = fn - if hasattr(self, "libby"): self._register_services({key: fn}) + if self.libby is not None: + self._register_services({key: fn}) - def add_services(self, mapping: Dict[str, RPCHandler]) -> None: - self.services.update(mapping) - if hasattr(self, "libby"): self._register_services(mapping) + def add_services(self, mapping: Mapping[str, RPCHandler]) -> None: + copied = dict(mapping) + self.services.update(copied) + if self.libby is not None: + self._register_services(copied) def add_topic(self, topic: str, fn: EvtHandler) -> None: self.topics[topic] = fn - if hasattr(self, "libby"): - self.libby.listen(topic, lambda msg, _h=fn: _h(msg.env.payload)) - self.libby.subscribe(topic) - - def add_topics(self, mapping: Dict[str, EvtHandler]) -> None: - self.topics.update(mapping) - if hasattr(self, "libby"): - for topic, fn in mapping.items(): - self.libby.listen(topic, lambda msg, _h=fn: _h(msg.env.payload)) - self.libby.subscribe(*mapping.keys()) - - # internals - def _must(self, name: str): - raise NotImplementedError(f"Set `{name}` or override config_{name}()") - - def _service_adapter(self, fn): - def adapter(user_payload: dict, _ctx: dict) -> dict: - try: - result = fn(user_payload) # user returns ANYTHING - return self.payload(result) # we "shove it into payload" for them - except Exception as ex: - return {"ok": False, "error": str(ex)} - return adapter + if self.libby is not None: + self._register_topics({topic: fn}) + + def add_topics(self, mapping: Mapping[str, EvtHandler]) -> None: + copied = dict(mapping) + self.topics.update(copied) + if self.libby is not None: + self._register_topics(copied) + + def register_keyword(self, keyword: Keyword) -> None: + """Register now, or defer registration until the daemon starts.""" + + if self.libby is None: + self._pending_keywords.append(keyword) + return + self.libby.register_keyword(keyword) + + def register_keywords(self, keywords: Iterable[Keyword]) -> None: + """Register many keywords now, or defer until startup.""" + + copied = list(keywords) + if self.libby is None: + self._pending_keywords.extend(copied) + return + self.libby.register_keywords(copied) + + @property + def keyword_registry(self): + """Return Libby's typed keyword builder. + + This property is intended for use in ``on_start``, after the Libby + instance has been constructed. + """ + + if self.libby is None: + raise RuntimeError( + "keyword_registry is available after daemon startup; " + "use register_keyword() to queue a keyword earlier" + ) + return self.libby.keyword_registry - def _register_services(self, mapping: Dict[str, RPCHandler]) -> None: - for key, fn in mapping.items(): - self.libby.serve_keys([key], self._service_adapter(fn)) + ### Libby construction def build_libby(self) -> Libby: - """Build Libby instance with selected transport.""" - if self.transport == "rabbitmq": + """Build a Libby instance for the configured transport.""" + + transport = str(self.transport).strip().lower() + + if transport == "rabbitmq": return Libby.rabbitmq( self_id=self.config_peer_id(), rabbitmq_url=self.config_rabbitmq_url(), @@ -164,66 +324,175 @@ def build_libby(self) -> Libby: callback=None, group_id=self.config_group_id(), ) - else: - # Default to ZMQ + + if transport == "zmq": return Libby.zmq( self_id=self.config_peer_id(), bind=self.config_bind(), address_book=self.config_address_book(), - keys=[], callback=None, # register per-key + keys=[], + callback=None, discover=self.config_discovery_enabled(), discover_interval_s=self.config_discovery_interval_s(), hello_on_start=True, group_id=self.config_group_id(), ) - def serve(self) -> None: - stop_evt = threading.Event() - def _sig(_s, _f): stop_evt.set() - signal.signal(signal.SIGINT, _sig) - signal.signal(signal.SIGTERM, _sig) + raise ValueError( + f"unsupported Libby transport {self.transport!r}; " + "expected 'zmq' or 'rabbitmq'" + ) + def start(self) -> None: + """Start Libby and initialize the daemon without blocking.""" + + if self._started: + return + + self._stop_event.clear() try: self.libby = self.build_libby() - except Exception as ex: - print(f"[{self.__class__.__name__}] failed to start: {ex}", file=sys.stderr) - raise - - if self.services: self._register_services(self.services) - if self.topics: - for topic, fn in self.topics.items(): - self.libby.listen(topic, lambda msg, _h=fn: _h(msg.env.payload)) - self.libby.subscribe(*self.topics.keys()) + self._register_topics(self.topics) - # discovery hello + hooks - try: if self.config_discovery_enabled(): - self.libby.hello() - self.on_hello(self.libby) + try: + self.libby.hello() + self.on_hello(self.libby) + except Exception: + self.logger.exception("discovery hello failed") + + try: + self.on_start(self.libby) + except Exception: + self.logger.exception("daemon initialization failed") + if self.fail_fast_on_start: + self._close_libby() + raise + + self._flush_keywords() + self._started = True + self.logger.info( + "started peer=%s transport=%s", + self.config_peer_id(), + self.transport, + ) except Exception: - pass + self._started = False + raise + + def request_stop(self) -> None: + """Request termination of a blocking ``serve`` call.""" + + self._stop_event.set() + + def stop(self) -> None: + """Stop the daemon; safe to call more than once.""" + + if self.libby is None and not self._started: + return try: - self.on_start(self.libby) - except Exception as ex: - print(f"[{self.__class__.__name__}] on_start error: {ex}", file=sys.stderr) + self.on_stop(self.libby) + except Exception: + self.logger.exception("daemon shutdown hook failed") + finally: + self._close_libby() + self._started = False + self._stop_event.set() + self.logger.info("stopped") + + def serve(self) -> None: + """Start the daemon and block until a signal or stop request.""" + + self._install_signal_handlers() + self.start() - if self.transport == "rabbitmq": - print(f"[{self.__class__.__name__}] up: id={self.config_peer_id()} transport=rabbitmq url={self.rabbitmq_url}") - else: - print(f"[{self.__class__.__name__}] up: id={self.config_peer_id()} bind={self.config_bind()}") try: - while not stop_evt.is_set(): time.sleep(0.5) + while not self._stop_event.wait(0.5): + pass finally: - try: self.on_stop() - except Exception: pass - self.libby.stop() - print(f"[{self.__class__.__name__}] stopped") + self.stop() + + ### Internals + + def _install_signal_handlers(self) -> None: + # signal.signal() is only legal in Python's main thread. Skipping it + # makes start/serve easier to exercise in test harnesses. + if threading.current_thread() is not threading.main_thread(): + self.logger.debug("signal handlers skipped outside main thread") + return + + def handle_signal(_signum: int, _frame: Any) -> None: + self.request_stop() + + signal.signal(signal.SIGINT, handle_signal) + signal.signal(signal.SIGTERM, handle_signal) + + def _must(self, name: str) -> Any: + raise ValueError( + f"set {name!r} in the class or configuration, " + f"or override config_{name}()" + ) + + def _service_adapter(self, fn: RPCHandler): + def adapter(user_payload: dict, _ctx: dict) -> dict: + try: + return self.payload(fn(user_payload)) + except Exception as exc: + self.logger.exception("service handler failed") + return {"ok": False, "error": str(exc)} + + return adapter + + def _register_services( + self, + mapping: Mapping[str, RPCHandler], + ) -> None: + if self.libby is None: + return + + for key, fn in mapping.items(): + self.libby.serve_keys([key], self._service_adapter(fn)) + + def _register_topics( + self, + mapping: Mapping[str, EvtHandler], + ) -> None: + if self.libby is None or not mapping: + return + + for topic, fn in mapping.items(): + self.libby.listen( + topic, + lambda msg, _handler=fn: _handler(msg.env.payload), + ) + self.libby.subscribe(*mapping.keys()) + + def _flush_keywords(self) -> None: + if self.libby is None: + return + + keywords = self._pending_keywords + self._pending_keywords = [] + keywords.extend(self.libby.keyword_registry.drain()) + + if keywords: + self.libby.register_keywords(keywords) + + def _close_libby(self) -> None: + libby, self.libby = self.libby, None + if libby is not None: + try: + libby.stop() + except Exception: + self.logger.exception("Libby transport shutdown failed") + + def payload(self, value: Any = None, /, **extra: Any) -> dict: + """Normalize a user result into a JSON-serializable dictionary.""" - def payload(self, value=None, /, **extra) -> dict: if value is None: - out = {} + out: Dict[str, Any] = {} elif is_dataclass(value): out = asdict(value) elif isinstance(value, cabc.Mapping): @@ -236,7 +505,9 @@ def payload(self, value=None, /, **extra) -> dict: try: json.dumps(out) - except TypeError as e: - raise ValueError(f"Payload not JSON-serializable: {e}") from e + except TypeError as exc: + raise ValueError( + f"payload is not JSON-serializable: {exc}" + ) from exc return out diff --git a/peers/peer_a.py b/peers/peer_a.py index d747b7e..6b40cca 100644 --- a/peers/peer_a.py +++ b/peers/peer_a.py @@ -1,10 +1,11 @@ import time + from libby.daemon import LibbyDaemon + class PeerA(LibbyDaemon): peer_id = "peer-A" - # ZMQ config (used when transport="zmq", which is the default) bind = "tcp://*:5555" address_book = { "peer-B": "tcp://127.0.0.1:5556", @@ -12,26 +13,44 @@ class PeerA(LibbyDaemon): "cli": "tcp://127.0.0.1:56001", } - # Transport selection: "zmq" (default) or "rabbitmq" transport = "zmq" - discovery_enabled = True discovery_interval_s = 2.0 def on_start(self, libby): try: - if not libby.wait_for_key("peer-B", "perf.echo", timeout_s=2.5): - libby.learn_peer_keys("peer-B", ["perf.echo", "ping.txt", "answer"]) - except AttributeError: - pass - - print("[PeerA] asking B: perf.echo …") - res = libby.rpc("peer-B", "perf.echo", {"t0": time.time()}, ttl_ms=8000) - print("[PeerA] result:", res) - - # publish a status - libby.publish("alerts.status", {"source": "peer-A", "ok": True}) + if not libby.wait_for_key( + "peer-B", + "perf.echo", + timeout_s=2.5, + ): + libby.learn_peer_keys( + "peer-B", + ["perf.echo", "ping.txt", "answer"], + ) + + print("[PeerA] asking B: perf.echo ...") + result = libby.rpc( + "peer-B", + "perf.echo", + {"t0": time.time()}, + ttl_ms=8000, + ) + print("[PeerA] result:", result) + + except Exception as exc: + print(f"[PeerA] Peer B request failed: {exc}") + + libby.publish( + "alerts.status", + { + "source": self.peer_id, + "ok": True, + "timestamp": time.time(), + }, + ) print("[PeerA] published alerts.status") + if __name__ == "__main__": PeerA().serve() diff --git a/peers/peer_b.py b/peers/peer_b.py index 97e2d05..5be05e3 100644 --- a/peers/peer_b.py +++ b/peers/peer_b.py @@ -1,50 +1,42 @@ import time -from typing import Dict, Any -from libby.daemon import LibbyDaemon - -def handle_echo(p: Dict[str, Any]): - # can return any JSON-serializable thing - return {"ok": True, "t0": p.get("t0"), "t1": time.time()} -def handle_ping(_p): - # returns a string - return "pong" - -def handle_answer(_p): - # returns a number - return 42 +from libby.daemon import LibbyDaemon -def on_status(payload: Dict[str, Any]) -> None: - print("[PeerB] alerts.status:", payload) class PeerB(LibbyDaemon): peer_id = "peer-B" - # ZMQ config (used when transport="zmq", which is the default) bind = "tcp://*:5556" address_book = { "peer-A": "tcp://127.0.0.1:5555", "peer-C": "tcp://127.0.0.1:5557", - "peer-D": "tcp://127.0.0.1:5558", "cli": "tcp://127.0.0.1:56001", } - # Transport selection: "zmq" (default) or "rabbitmq" transport = "zmq" - discovery_enabled = True discovery_interval_s = 2.0 services = { - "perf.echo": handle_echo, - "ping.txt": handle_ping, - "answer": handle_answer, + "perf.echo": lambda payload: { + "ok": True, + "received": payload, + "responded_at": time.time(), + }, + "ping.txt": lambda payload: { + "ok": True, + "message": "pong", + }, + "answer": lambda payload: { + "ok": True, + "value": 42, + }, } - # Topics (PUB/SUB) - topics = { - "alerts.status": on_status, - } + def on_start(self, libby): + print("[PeerB] ready") + print("[PeerB] services:", list(self.services)) + if __name__ == "__main__": PeerB().serve() diff --git a/peers/peer_c.py b/peers/peer_c.py index 8b771b1..24186da 100644 --- a/peers/peer_c.py +++ b/peers/peer_c.py @@ -1,52 +1,40 @@ -import time -from typing import Dict, Any from libby.daemon import LibbyDaemon -def info(_p: Dict[str, Any]): - return {"ok": True, "info": "peer-C", "time": time.time()} - -def math_add(p: Dict[str, Any]): - a, b = p.get("a"), p.get("b") - if not isinstance(a, (int, float)) or not isinstance(b, (int, float)): - return {"ok": False, "error": "need numeric a and b"} - return {"ok": True, "sum": a + b} class PeerC(LibbyDaemon): peer_id = "peer-C" - # ZMQ config (used when transport="zmq", which is the default) bind = "tcp://*:5557" address_book = { "peer-A": "tcp://127.0.0.1:5555", "peer-B": "tcp://127.0.0.1:5556", - "peer-C": "tcp://127.0.0.1:5557" + "cli": "tcp://127.0.0.1:56001", } - # Transport selection: "zmq" (default) or "rabbitmq" transport = "zmq" - rabbitmq_url = "amqp://localhost" # Used when transport="rabbitmq" - discovery_enabled = True discovery_interval_s = 2.0 - services = { - "clientC.info": info, - "math.add": math_add, - } + def __init__(self): + super().__init__() + + self.last_status = None + self.add_topic("alerts.status", self.handle_status) + self.add_service("alerts.last", self.get_last_status) + + def handle_status(self, payload): + self.last_status = payload + print("[PeerC] received alerts.status:", payload) + + def get_last_status(self, payload): + return { + "ok": True, + "status": self.last_status, + } def on_start(self, libby): - # Proxy service - def echo_proxy(_p: Dict[str, Any]): - res = libby.rpc("peer-B", "perf.echo", {"t0": time.time()}, ttl_ms=6000) - # returning a dict - return {"ok": True, "forwarded_to": "peer-B", "result": res} - - # Register at runtime - self.add_service("perf.echo.proxy", echo_proxy) - - print("[PeerC] math.add(2,5) ->", - libby.rpc(self.peer_id, "math.add", {"a": 2, "b": 5}, ttl_ms=2000)) - libby.publish("alerts.status", {"source": "peer-C", "ok": True}) + print("[PeerC] listening for alerts.status") + if __name__ == "__main__": PeerC().serve()