From 37e94bcb6cd2f97fed4fe8bfc7275aeae970b67f Mon Sep 17 00:00:00 2001 From: gogocat Date: Tue, 7 Apr 2026 14:25:39 +0300 Subject: [PATCH 1/2] =?UTF-8?q?chore:=20bump=20version=200.14.19=20?= =?UTF-8?q?=E2=86=92=200.14.21,=20update=20changelog=20and=20roadmap?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Release v0.14.21: - Redis Cacher production-ready (start/stop lifecycle, 16 integration tests) - seq/instanceID heartbeat checks (Node.js parity) - 10 audit findings fixed (3 CRITICAL: instanceID persistence, payload guard, logger fallback + 4 HIGH + 3 MEDIUM) Evidence: 2374 tests, 28/28 demo matrix, 90/90 comprehensive, mypy 0. Co-Authored-By: Claude Opus 4.6 (1M context) --- CHANGELOG.md | 22 ++++++++++++++++++++++ moleculerpy/__init__.py | 2 +- pyproject.toml | 2 +- 3 files changed, 24 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2fcf17e..d84139c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,28 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [0.14.21] - 2026-04-07 + +### Added +- **Redis Cacher production-ready** — start()/stop() lifecycle, broker integration, + 16 integration tests (get/set/delete/clean/TTL/keys/broker) +- **seq/instanceID heartbeat checks** — detects remote node restart and service changes + via heartbeat (Node.js heartbeatReceived parity) + +### Fixed +- **instanceID persistence** — process_node_info now saves instanceID from INFO packets + (was causing infinite re-discovery loop on every heartbeat) +- **payload dict guard** — _handle_heartbeat validates payload is dict before access + (prevents crash on malformed packets) +- **seq type coercion** — int comparison for cross-language safety +- **instanceID=None guard** — skip comparison when node has no instanceID yet +- **_discover_pending cleanup** — stale entries evicted in check_remote_nodes + (prevents unbounded memory growth) +- **Redis cacher: logger fallback** — set in __init__ (was crashing if connect() before init()) +- **Redis cacher: prefix dedup** — removed duplicate namespace logic (BaseCacher handles it) +- **Redis cacher: TTL validation** — negative/zero TTL warns and stores without expiry +- **Redis cacher: ping loop** — pings first then sleeps (detects immediate connection drop) + ## [0.14.19] - 2026-04-06 ### Fixed (Audit-driven fixes for v0.14.18 serializers — PRD-019) diff --git a/moleculerpy/__init__.py b/moleculerpy/__init__.py index c3c4471..988acf2 100644 --- a/moleculerpy/__init__.py +++ b/moleculerpy/__init__.py @@ -54,7 +54,7 @@ try: __version__ = version("moleculerpy") except PackageNotFoundError: - __version__ = "0.14.19" + __version__ = "0.14.21" __all__ = [ # noqa: RUF022 # Core diff --git a/pyproject.toml b/pyproject.toml index 228dc41..ba07571 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "moleculerpy" -version = "0.14.19" +version = "0.14.21" description = "Fast, modern microservices framework for Python - Port of Moleculer.js" authors = [ { name = "Eli Rum", email = "explosivebit@gmail.com" } From ce04f31256a4d4f0f4dd549c47f5fab3459db02f Mon Sep 17 00:00:00 2001 From: gogocat Date: Tue, 7 Apr 2026 15:54:54 +0300 Subject: [PATCH 2/2] =?UTF-8?q?feat(protocol,lifecycle):=20protocol=20corr?= =?UTF-8?q?ectness=20sprint=20=E2=80=94=206=20tasks=20in=20parallel?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Sprint executed via TeamCreate with 6 teammates coordinated by team-lead. Node.js parity fixes for protocol + graceful lifecycle. ## Tasks completed 1. **Broker hook rename with aliases** (hook-renamer) - middleware/base.py: added short aliases starting/started/stopping - broker.py: _call_middleware_hooks invokes both long + short names - Backward compatible — existing broker_* middleware works unchanged - Skipped `stopped` alias due to collision with existing no-arg cleanup hook 2. **PacketHeartbeat proto schema extension** (proto-extender) - serializers/proto/packets.proto: added fields 4-7 (seq, instanceID, memory, cpuSeq) - Regenerated packets_pb2.py - protobuf.py: _HEARTBEAT_MAX_FIELDS 4 → 8 - ADR-heartbeat-schema.md documenting the decision - Wire-compatible with Node.js (ignores unknown fields) 3. **TrackingConfig in Settings** (tracking-configurator) - settings.py: TrackingConfig dataclass (enabled=False, shutdown_timeout=5.0s) - Validation in Settings._validate() - Exported via moleculerpy/__init__.py 4. **Connection drain on broker stop** (connection-drainer) - transit.py: send_disconnect_info() broadcasts INFO(services=[]) before disconnect - broker.stop(): calls send_disconnect_info() BEFORE transit.disconnect() - Matches Node.js service-broker.js:531-539 pattern - e2e test with 2 memory brokers verifies remote node drops math.add before DISCONNECT 5. **ContextTracker auto-registration** (tracker-integrator) - broker.py __init__: auto-append ContextTrackerMiddleware if settings.tracking.enabled - Converts float seconds → int ms for middleware - e2e test: slow action in-flight, stop() waits for completion 6. **Protocol lifecycle test suite** (test-author) - tests/unit/protocol_lifecycle_test.py (12 tests, ~270 LOC) - System-level integration tests for all 5 changes above ## Evidence - ruff format + check: clean - mypy --strict: 0 errors - pytest: **2396 passed** (2378 baseline + 18 new) - demo_matrix: 28/28 OK - demo_comprehensive: 90/90 OK Refs: Protocol gap audit report Co-Authored-By: Claude Opus 4.6 (1M context) --- moleculerpy/__init__.py | 3 +- moleculerpy/broker.py | 62 +++- moleculerpy/middleware/base.py | 21 ++ moleculerpy/serializers/proto/packets.proto | 14 +- moleculerpy/serializers/proto/packets_pb2.py | 28 +- moleculerpy/serializers/protobuf.py | 5 +- moleculerpy/settings.py | 28 ++ moleculerpy/transit.py | 18 ++ tests/e2e/test_drain_on_stop.py | 63 ++++ tests/e2e/test_tracking_shutdown.py | 58 ++++ tests/unit/broker_test.py | 22 ++ tests/unit/protocol_lifecycle_test.py | 296 +++++++++++++++++++ tests/unit/serializer_cbor_protobuf_test.py | 20 ++ tests/unit/settings_test.py | 24 +- tests/unit/transit_test.py | 45 +++ 15 files changed, 678 insertions(+), 29 deletions(-) create mode 100644 tests/e2e/test_drain_on_stop.py create mode 100644 tests/e2e/test_tracking_shutdown.py create mode 100644 tests/unit/protocol_lifecycle_test.py diff --git a/moleculerpy/__init__.py b/moleculerpy/__init__.py index 988acf2..fb03a81 100644 --- a/moleculerpy/__init__.py +++ b/moleculerpy/__init__.py @@ -39,7 +39,7 @@ ) from .serializers import BaseSerializer, JsonSerializer, MsgPackSerializer, resolve_serializer from .service import Service -from .settings import Settings, SettingsValidationError +from .settings import Settings, SettingsValidationError, TrackingConfig from .stream import AsyncStream, StreamError from .tracing import ( BaseTraceExporter, @@ -65,6 +65,7 @@ "Lifecycle", "Settings", "SettingsValidationError", + "TrackingConfig", "NodeID", "ServiceName", "ActionName", diff --git a/moleculerpy/broker.py b/moleculerpy/broker.py index abf97c3..3a0a114 100644 --- a/moleculerpy/broker.py +++ b/moleculerpy/broker.py @@ -137,6 +137,16 @@ def __init__( self._validator = resolve_validator(getattr(self.settings, "validator", "default")) + # Auto-register ContextTracker middleware if tracking enabled + tracking_cfg = getattr(self.settings, "tracking", None) + if tracking_cfg is not None and getattr(tracking_cfg, "enabled", False): + from .middleware.context_tracker import ContextTrackerMiddleware # noqa: PLC0415 + + # TrackingConfig.shutdown_timeout is float seconds; + # ContextTrackerMiddleware expects int milliseconds. + shutdown_timeout_ms = int(tracking_cfg.shutdown_timeout * 1000) + self.middlewares.append(ContextTrackerMiddleware(shutdown_timeout=shutdown_timeout_ms)) + # Wrapped event methods (set during start() by middleware) self._wrapped_emit: ( Callable[[str, dict[str, Any], dict[str, Any]], Awaitable[Any]] | None @@ -264,21 +274,50 @@ def _call_middleware_hooks( Returns: List of coroutines if is_async=True, None otherwise """ + # Node.js Moleculer uses short hook names (starting/started/stopping/stopped) + # while MoleculerPy historically used broker_* names. To maintain backward + # compatibility AND Node.js ecosystem compatibility, we invoke both names. + # Note: "stopped" is intentionally NOT aliased because MoleculerPy's existing + # stopped() hook takes no arguments (middleware self-cleanup), which would + # collide with Node.js stopped(broker) signature. + from .middleware.base import Middleware as _BaseMiddleware # noqa: PLC0415 + + _broker_hook_aliases = { + "broker_starting": "starting", + "broker_started": "started", + "broker_stopping": "stopping", + } + alias = _broker_hook_aliases.get(hook_name) + + def _is_overridden(mw: Any, name: str) -> bool: + """True if mw class overrides the alias method (not base no-op).""" + mw_method = getattr(type(mw), name, None) + base_method = getattr(_BaseMiddleware, name, None) + return mw_method is not None and mw_method is not base_method + if is_async: coroutines = [] for middleware in self.middlewares: - hook = getattr(middleware, hook_name, None) - if hook and callable(hook): - result = hook(*args) - if asyncio.iscoroutine(result): - coroutines.append(result) + names = [hook_name] + if alias and _is_overridden(middleware, alias): + names.append(alias) + for name in names: + hook = getattr(middleware, name, None) + if hook and callable(hook): + result = hook(*args) + if asyncio.iscoroutine(result): + coroutines.append(result) return coroutines if coroutines else None else: # Synchronous hooks for middleware in self.middlewares: - hook = getattr(middleware, hook_name, None) - if hook and callable(hook): - hook(*args) + names = [hook_name] + if alias and _is_overridden(middleware, alias): + names.append(alias) + for name in names: + hook = getattr(middleware, name, None) + if hook and callable(hook): + hook(*args) return None async def _execute_middleware_hooks( @@ -559,6 +598,13 @@ async def _stop_core() -> None: if self.cacher: await self.cacher.stop() + # Drain: notify peers we're shutting down (empty services) so + # they stop routing new requests to us BEFORE we DISCONNECT. + try: + await self.transit.send_disconnect_info() + except Exception as e: + self.logger.warning(f"Error sending drain INFO: {e}") + # Disconnect from the cluster await self.transit.disconnect() diff --git a/moleculerpy/middleware/base.py b/moleculerpy/middleware/base.py index 6f902f5..b7f180c 100644 --- a/moleculerpy/middleware/base.py +++ b/moleculerpy/middleware/base.py @@ -156,6 +156,27 @@ async def broker_stopped(self, broker: Any) -> None: """ pass + # ========================================================================== + # Node.js Moleculer-compatible short aliases for broker lifecycle hooks. + # The broker invokes BOTH the long-form (broker_*) and the short alias name + # so middleware authored against the Node.js naming convention works as-is. + # Note: "stopped" is intentionally NOT aliased because MoleculerPy's existing + # stopped() hook (below) takes no arguments and is reserved for middleware + # self-cleanup. Use broker_stopped() if you need a broker reference. + # ========================================================================== + + async def starting(self, broker: Any) -> None: + """Node.js-compatible alias for broker_starting. Called BEFORE connect.""" + pass + + async def started(self, broker: Any) -> None: + """Node.js-compatible alias for broker_started. Called AFTER connect.""" + pass + + async def stopping(self, broker: Any) -> None: + """Node.js-compatible alias for broker_stopping. Called BEFORE disconnect.""" + pass + async def service_creating(self, service: Any) -> None: """ Hook called asynchronously BEFORE a service is registered. diff --git a/moleculerpy/serializers/proto/packets.proto b/moleculerpy/serializers/proto/packets.proto index e2281e1..4bc70fc 100644 --- a/moleculerpy/serializers/proto/packets.proto +++ b/moleculerpy/serializers/proto/packets.proto @@ -92,9 +92,17 @@ message PacketDisconnect { } message PacketHeartbeat { - string ver = 1; - string sender = 2; - double cpu = 3; + string ver = 1; + string sender = 2; + double cpu = 3; + // MoleculerPy extension fields (4-7) for restart detection. + // Wire-compatible with Node.js: unknown fields are silently ignored + // by Node.js ProtoBuf parser. Field numbers 4-7 are RESERVED FOREVER. + // See: .forgeplan/adrs/ADR-heartbeat-schema.md + int32 seq = 4; + string instanceID = 5; + double memory = 6; + int32 cpuSeq = 7; } message PacketPing { diff --git a/moleculerpy/serializers/proto/packets_pb2.py b/moleculerpy/serializers/proto/packets_pb2.py index fe1f4d2..4a8d0ae 100644 --- a/moleculerpy/serializers/proto/packets_pb2.py +++ b/moleculerpy/serializers/proto/packets_pb2.py @@ -19,7 +19,7 @@ DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( - b'\n\rpackets.proto\x12\x07packets"\xac\x02\n\x0bPacketEvent\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\r\n\x05\x65vent\x18\x04 \x01(\t\x12\x0c\n\x04\x64\x61ta\x18\x05 \x01(\x0c\x12#\n\x08\x64\x61taType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\x0e\n\x06groups\x18\x07 \x03(\t\x12\x0c\n\x04meta\x18\t \x01(\t\x12\x11\n\tbroadcast\x18\x08 \x01(\x08\x12\r\n\x05level\x18\n \x01(\x05\x12\x0f\n\x07tracing\x18\x0b \x01(\x08\x12\x10\n\x08parentID\x18\x0c \x01(\t\x12\x11\n\trequestID\x18\r \x01(\t\x12\x0e\n\x06stream\x18\x0e \x01(\x08\x12\x0b\n\x03seq\x18\x0f \x01(\x05\x12\x0e\n\x06\x63\x61ller\x18\x10 \x01(\t\x12\x0f\n\x07needAck\x18\x11 \x01(\x08"\x90\x02\n\rPacketRequest\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\x0e\n\x06\x61\x63tion\x18\x04 \x01(\t\x12\x0e\n\x06params\x18\x05 \x01(\x0c\x12%\n\nparamsType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\x0c\n\x04meta\x18\x07 \x01(\t\x12\x0f\n\x07timeout\x18\x08 \x01(\x01\x12\r\n\x05level\x18\t \x01(\x05\x12\x0f\n\x07tracing\x18\n \x01(\x08\x12\x10\n\x08parentID\x18\x0b \x01(\t\x12\x11\n\trequestID\x18\x0c \x01(\t\x12\x0e\n\x06stream\x18\r \x01(\x08\x12\x0b\n\x03seq\x18\x0e \x01(\x05\x12\x0e\n\x06\x63\x61ller\x18\x0f \x01(\t"\xb7\x01\n\x0ePacketResponse\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\x0f\n\x07success\x18\x04 \x01(\x08\x12\x0c\n\x04\x64\x61ta\x18\x05 \x01(\x0c\x12#\n\x08\x64\x61taType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\r\n\x05\x65rror\x18\x07 \x01(\t\x12\x0c\n\x04meta\x18\x08 \x01(\t\x12\x0e\n\x06stream\x18\t \x01(\x08\x12\x0b\n\x03seq\x18\n \x01(\x05"-\n\x0ePacketDiscover\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t"\x8a\x02\n\nPacketInfo\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x10\n\x08services\x18\x03 \x01(\t\x12\x0e\n\x06\x63onfig\x18\x04 \x01(\t\x12\x0e\n\x06ipList\x18\x05 \x03(\t\x12\x10\n\x08hostname\x18\x06 \x01(\t\x12*\n\x06\x63lient\x18\x07 \x01(\x0b\x32\x1a.packets.PacketInfo.Client\x12\x0b\n\x03seq\x18\x08 \x01(\x05\x12\x12\n\ninstanceID\x18\t \x01(\t\x12\x10\n\x08metadata\x18\n \x01(\t\x1a<\n\x06\x43lient\x12\x0c\n\x04type\x18\x01 \x01(\t\x12\x0f\n\x07version\x18\x02 \x01(\t\x12\x13\n\x0blangVersion\x18\x03 \x01(\t"/\n\x10PacketDisconnect\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t";\n\x0fPacketHeartbeat\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0b\n\x03\x63pu\x18\x03 \x01(\x01"C\n\nPacketPing\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04time\x18\x03 \x01(\x03\x12\n\n\x02id\x18\x04 \x01(\t"T\n\nPacketPong\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04time\x18\x03 \x01(\x03\x12\x0f\n\x07\x61rrived\x18\x04 \x01(\x03\x12\n\n\x02id\x18\x05 \x01(\t"L\n\x11PacketGossipHello\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04host\x18\x03 \x01(\t\x12\x0c\n\x04port\x18\x04 \x01(\x05"S\n\x13PacketGossipRequest\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0e\n\x06online\x18\x03 \x01(\t\x12\x0f\n\x07offline\x18\x04 \x01(\t"T\n\x14PacketGossipResponse\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0e\n\x06online\x18\x03 \x01(\t\x12\x0f\n\x07offline\x18\x04 \x01(\t*]\n\x08\x44\x61taType\x12\x16\n\x12\x44\x41TATYPE_UNDEFINED\x10\x00\x12\x11\n\rDATATYPE_NULL\x10\x01\x12\x11\n\rDATATYPE_JSON\x10\x02\x12\x13\n\x0f\x44\x41TATYPE_BUFFER\x10\x03\x62\x06proto3' + b'\n\rpackets.proto\x12\x07packets"\xac\x02\n\x0bPacketEvent\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\r\n\x05\x65vent\x18\x04 \x01(\t\x12\x0c\n\x04\x64\x61ta\x18\x05 \x01(\x0c\x12#\n\x08\x64\x61taType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\x0e\n\x06groups\x18\x07 \x03(\t\x12\x0c\n\x04meta\x18\t \x01(\t\x12\x11\n\tbroadcast\x18\x08 \x01(\x08\x12\r\n\x05level\x18\n \x01(\x05\x12\x0f\n\x07tracing\x18\x0b \x01(\x08\x12\x10\n\x08parentID\x18\x0c \x01(\t\x12\x11\n\trequestID\x18\r \x01(\t\x12\x0e\n\x06stream\x18\x0e \x01(\x08\x12\x0b\n\x03seq\x18\x0f \x01(\x05\x12\x0e\n\x06\x63\x61ller\x18\x10 \x01(\t\x12\x0f\n\x07needAck\x18\x11 \x01(\x08"\x90\x02\n\rPacketRequest\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\x0e\n\x06\x61\x63tion\x18\x04 \x01(\t\x12\x0e\n\x06params\x18\x05 \x01(\x0c\x12%\n\nparamsType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\x0c\n\x04meta\x18\x07 \x01(\t\x12\x0f\n\x07timeout\x18\x08 \x01(\x01\x12\r\n\x05level\x18\t \x01(\x05\x12\x0f\n\x07tracing\x18\n \x01(\x08\x12\x10\n\x08parentID\x18\x0b \x01(\t\x12\x11\n\trequestID\x18\x0c \x01(\t\x12\x0e\n\x06stream\x18\r \x01(\x08\x12\x0b\n\x03seq\x18\x0e \x01(\x05\x12\x0e\n\x06\x63\x61ller\x18\x0f \x01(\t"\xb7\x01\n\x0ePacketResponse\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\n\n\x02id\x18\x03 \x01(\t\x12\x0f\n\x07success\x18\x04 \x01(\x08\x12\x0c\n\x04\x64\x61ta\x18\x05 \x01(\x0c\x12#\n\x08\x64\x61taType\x18\x06 \x01(\x0e\x32\x11.packets.DataType\x12\r\n\x05\x65rror\x18\x07 \x01(\t\x12\x0c\n\x04meta\x18\x08 \x01(\t\x12\x0e\n\x06stream\x18\t \x01(\x08\x12\x0b\n\x03seq\x18\n \x01(\x05"-\n\x0ePacketDiscover\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t"\x8a\x02\n\nPacketInfo\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x10\n\x08services\x18\x03 \x01(\t\x12\x0e\n\x06\x63onfig\x18\x04 \x01(\t\x12\x0e\n\x06ipList\x18\x05 \x03(\t\x12\x10\n\x08hostname\x18\x06 \x01(\t\x12*\n\x06\x63lient\x18\x07 \x01(\x0b\x32\x1a.packets.PacketInfo.Client\x12\x0b\n\x03seq\x18\x08 \x01(\x05\x12\x12\n\ninstanceID\x18\t \x01(\t\x12\x10\n\x08metadata\x18\n \x01(\t\x1a<\n\x06\x43lient\x12\x0c\n\x04type\x18\x01 \x01(\t\x12\x0f\n\x07version\x18\x02 \x01(\t\x12\x13\n\x0blangVersion\x18\x03 \x01(\t"/\n\x10PacketDisconnect\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t"|\n\x0fPacketHeartbeat\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0b\n\x03\x63pu\x18\x03 \x01(\x01\x12\x0b\n\x03seq\x18\x04 \x01(\x05\x12\x12\n\ninstanceID\x18\x05 \x01(\t\x12\x0e\n\x06memory\x18\x06 \x01(\x01\x12\x0e\n\x06\x63puSeq\x18\x07 \x01(\x05"C\n\nPacketPing\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04time\x18\x03 \x01(\x03\x12\n\n\x02id\x18\x04 \x01(\t"T\n\nPacketPong\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04time\x18\x03 \x01(\x03\x12\x0f\n\x07\x61rrived\x18\x04 \x01(\x03\x12\n\n\x02id\x18\x05 \x01(\t"L\n\x11PacketGossipHello\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0c\n\x04host\x18\x03 \x01(\t\x12\x0c\n\x04port\x18\x04 \x01(\x05"S\n\x13PacketGossipRequest\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0e\n\x06online\x18\x03 \x01(\t\x12\x0f\n\x07offline\x18\x04 \x01(\t"T\n\x14PacketGossipResponse\x12\x0b\n\x03ver\x18\x01 \x01(\t\x12\x0e\n\x06sender\x18\x02 \x01(\t\x12\x0e\n\x06online\x18\x03 \x01(\t\x12\x0f\n\x07offline\x18\x04 \x01(\t*]\n\x08\x44\x61taType\x12\x16\n\x12\x44\x41TATYPE_UNDEFINED\x10\x00\x12\x11\n\rDATATYPE_NULL\x10\x01\x12\x11\n\rDATATYPE_JSON\x10\x02\x12\x13\n\x0f\x44\x41TATYPE_BUFFER\x10\x03\x62\x06proto3' ) _globals = globals() @@ -27,8 +27,8 @@ _builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, "packets_pb2", _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None - _globals["_DATATYPE"]._serialized_start = 1620 - _globals["_DATATYPE"]._serialized_end = 1713 + _globals["_DATATYPE"]._serialized_start = 1685 + _globals["_DATATYPE"]._serialized_end = 1778 _globals["_PACKETEVENT"]._serialized_start = 27 _globals["_PACKETEVENT"]._serialized_end = 327 _globals["_PACKETREQUEST"]._serialized_start = 330 @@ -44,15 +44,15 @@ _globals["_PACKETDISCONNECT"]._serialized_start = 1106 _globals["_PACKETDISCONNECT"]._serialized_end = 1153 _globals["_PACKETHEARTBEAT"]._serialized_start = 1155 - _globals["_PACKETHEARTBEAT"]._serialized_end = 1214 - _globals["_PACKETPING"]._serialized_start = 1216 - _globals["_PACKETPING"]._serialized_end = 1283 - _globals["_PACKETPONG"]._serialized_start = 1285 - _globals["_PACKETPONG"]._serialized_end = 1369 - _globals["_PACKETGOSSIPHELLO"]._serialized_start = 1371 - _globals["_PACKETGOSSIPHELLO"]._serialized_end = 1447 - _globals["_PACKETGOSSIPREQUEST"]._serialized_start = 1449 - _globals["_PACKETGOSSIPREQUEST"]._serialized_end = 1532 - _globals["_PACKETGOSSIPRESPONSE"]._serialized_start = 1534 - _globals["_PACKETGOSSIPRESPONSE"]._serialized_end = 1618 + _globals["_PACKETHEARTBEAT"]._serialized_end = 1279 + _globals["_PACKETPING"]._serialized_start = 1281 + _globals["_PACKETPING"]._serialized_end = 1348 + _globals["_PACKETPONG"]._serialized_start = 1350 + _globals["_PACKETPONG"]._serialized_end = 1434 + _globals["_PACKETGOSSIPHELLO"]._serialized_start = 1436 + _globals["_PACKETGOSSIPHELLO"]._serialized_end = 1512 + _globals["_PACKETGOSSIPREQUEST"]._serialized_start = 1514 + _globals["_PACKETGOSSIPREQUEST"]._serialized_end = 1597 + _globals["_PACKETGOSSIPRESPONSE"]._serialized_start = 1599 + _globals["_PACKETGOSSIPRESPONSE"]._serialized_end = 1683 # @@protoc_insertion_point(module_scope) diff --git a/moleculerpy/serializers/protobuf.py b/moleculerpy/serializers/protobuf.py index 6cc4126..bccdfe8 100644 --- a/moleculerpy/serializers/protobuf.py +++ b/moleculerpy/serializers/protobuf.py @@ -47,8 +47,9 @@ # Separate from BaseSerializer.MAX_PAYLOAD_BYTES which guards the whole frame. MAX_NESTED_FIELD_BYTES: Final[int] = 1 * 1024 * 1024 # 1MB per nested field -# Heuristic constant: HEARTBEAT packet has small field count (ver, sender, cpu, +1 optional) -_HEARTBEAT_MAX_FIELDS: Final[int] = 4 +# Heuristic constant: HEARTBEAT packet field count. +# Schema: ver, sender, cpu, seq, instanceID, memory, cpuSeq (+1 optional slack). +_HEARTBEAT_MAX_FIELDS: Final[int] = 8 def _check_json_depth(text: str, max_depth: int = MAX_JSON_DEPTH) -> bool: diff --git a/moleculerpy/settings.py b/moleculerpy/settings.py index 740cf59..c73cb7a 100644 --- a/moleculerpy/settings.py +++ b/moleculerpy/settings.py @@ -1,3 +1,4 @@ +from dataclasses import dataclass from typing import TYPE_CHECKING, Any, ClassVar if TYPE_CHECKING: @@ -12,6 +13,25 @@ class SettingsValidationError(ValueError): pass +@dataclass +class TrackingConfig: + """Configuration for context tracking / graceful shutdown. + + Mirrors Node.js Moleculer's ``tracking`` broker option. + + Attributes: + enabled: If True, ContextTracker tracks active contexts so the + broker can wait for them to complete during graceful stop. + Defaults to False, matching Node.js Moleculer. + shutdown_timeout: Maximum time (in seconds) to wait for in-flight + contexts to finish during shutdown. Defaults to 5.0 seconds + (Node.js default is 5000ms). + """ + + enabled: bool = False + shutdown_timeout: float = 5.0 + + class Settings: """Configuration settings for the MoleculerPy broker. @@ -84,6 +104,7 @@ def __init__( namespace: str | None = None, disable_balancer: bool = False, validator: "str | bool | type[BaseValidator] | BaseValidator | None" = "default", + tracking: TrackingConfig | None = None, ) -> None: self.transporter = transporter self.serializer = serializer @@ -103,6 +124,7 @@ def __init__( self.namespace = namespace self.disable_balancer = disable_balancer self.validator = validator + self.tracking = tracking if tracking is not None else TrackingConfig() # Validate all settings self._validate() @@ -172,6 +194,12 @@ def _validate(self) -> None: f"got '{self.transporter}'" ) + # Validate tracking + if self.tracking.shutdown_timeout <= 0: + raise SettingsValidationError( + f"tracking.shutdown_timeout must be positive, got {self.tracking.shutdown_timeout}" + ) + # Validate strategy if self.strategy.upper() not in self.VALID_STRATEGIES: raise SettingsValidationError( diff --git a/moleculerpy/transit.py b/moleculerpy/transit.py index f8af5d3..41018b0 100644 --- a/moleculerpy/transit.py +++ b/moleculerpy/transit.py @@ -540,6 +540,24 @@ async def send_node_info(self) -> None: node_info = self.node_catalog.local_node.get_info() await self.publish(Packet(Topic.INFO, None, node_info)) + async def send_disconnect_info(self) -> None: + """Broadcast INFO packet with empty services list to drain connections. + + Sent BEFORE DISCONNECT during graceful shutdown so peer nodes mark this + node as draining and stop routing new requests to it. Matches Node.js + Moleculer service-broker.js stop() pattern. + """ + if not self._was_connected: + return + if self.node_catalog.local_node is None: + return + try: + info = self.node_catalog.local_node.get_info() + drain_info = {**info, "services": []} + await self.publish(Packet(Topic.INFO, None, drain_info)) + except Exception as e: + self.logger.warning(f"Error sending disconnect INFO drain: {e}") + async def _handle_discover(self, packet: Packet) -> None: """Handle discovery requests by sending node info. diff --git a/tests/e2e/test_drain_on_stop.py b/tests/e2e/test_drain_on_stop.py new file mode 100644 index 0000000..9b8adaf --- /dev/null +++ b/tests/e2e/test_drain_on_stop.py @@ -0,0 +1,63 @@ +"""E2E test: broker.stop() drains services on remote nodes before disconnecting. + +Verifies the Node.js Moleculer pattern: an INFO packet with empty services list +is broadcast BEFORE the DISCONNECT packet, so peer nodes mark the node as +draining and stop routing new requests to it. +""" + +from __future__ import annotations + +import asyncio + +import pytest + +from moleculerpy.broker import ServiceBroker +from moleculerpy.decorators import action +from moleculerpy.service import Service +from moleculerpy.settings import Settings + + +class MathDrainService(Service): + name = "math" + + def __init__(self) -> None: + super().__init__(self.name) + + @action() + async def add(self, ctx) -> int: + return ctx.params["a"] + ctx.params["b"] + + +@pytest.mark.asyncio +@pytest.mark.e2e +async def test_drain_info_sent_before_disconnect() -> None: + """Broker A stops → Broker B sees math service removed before DISCONNECT.""" + settings_a = Settings(transporter="memory://", prefer_local=False) + broker_a = ServiceBroker(id="node-a", settings=settings_a) + await broker_a.register(MathDrainService()) + + settings_b = Settings(transporter="memory://", prefer_local=False) + broker_b = ServiceBroker(id="node-b", settings=settings_b) + + await broker_a.start() + await broker_b.start() + + try: + # Wait until broker B sees math service from node-a + await broker_b.wait_for_services(["math"], timeout=5.0) + assert broker_b.registry.get_action("math.add") is not None + + # Stop broker A — this should trigger drain INFO before DISCONNECT + await broker_a.stop() + + # Allow propagation + await asyncio.sleep(0.2) + + # Broker B should no longer have math.add from node-a + action_obj = broker_b.registry.get_action("math.add") + # Either action is gone or the node-a entry was removed + assert action_obj is None or not any( + getattr(ep, "node_id", None) == "node-a" for ep in getattr(action_obj, "endpoints", []) + ) + finally: + await broker_b.stop() diff --git a/tests/e2e/test_tracking_shutdown.py b/tests/e2e/test_tracking_shutdown.py new file mode 100644 index 0000000..7f1f46f --- /dev/null +++ b/tests/e2e/test_tracking_shutdown.py @@ -0,0 +1,58 @@ +"""E2E test: ContextTracker auto-registration drains in-flight requests on stop. + +When `settings.tracking.enabled=True`, broker.stop() should wait for +in-flight contexts to complete (up to shutdown_timeout) before disconnecting. +""" + +from __future__ import annotations + +import asyncio + +import pytest + +from moleculerpy.broker import ServiceBroker +from moleculerpy.decorators import action +from moleculerpy.service import Service +from moleculerpy.settings import Settings, TrackingConfig + + +class SlowService(Service): + name = "slow" + + def __init__(self) -> None: + super().__init__(self.name) + self.completed = False + + @action() + async def work(self, ctx) -> str: + await asyncio.sleep(1.0) + self.completed = True + return "done" + + +@pytest.mark.asyncio +@pytest.mark.e2e +async def test_tracking_shutdown_waits_for_inflight() -> None: + """broker.stop() should wait for in-flight tracked action to complete.""" + settings = Settings( + transporter="memory://", + tracking=TrackingConfig(enabled=True, shutdown_timeout=5.0), + ) + broker = ServiceBroker(id="track-node", settings=settings) + svc = SlowService() + await broker.register(svc) + await broker.start() + + # Fire request in background + call_task = asyncio.create_task(broker.call("slow.work")) + # Let it begin + await asyncio.sleep(0.1) + assert not svc.completed + + # Stop should wait for the in-flight action to finish + await broker.stop() + + # Action should have completed before stop returned + assert svc.completed + assert call_task.done() + assert await call_task == "done" diff --git a/tests/unit/broker_test.py b/tests/unit/broker_test.py index 83b60e5..3e5a81e 100644 --- a/tests/unit/broker_test.py +++ b/tests/unit/broker_test.py @@ -362,3 +362,25 @@ async def test_broker_call_remote_action_with_error( await broker.call("remote.error") mock_transit.request.assert_called_once_with(endpoint, context) + + +def test_tracking_disabled_no_middleware(): + """Default settings — ContextTrackerMiddleware is NOT auto-registered.""" + from moleculerpy.middleware.context_tracker import ContextTrackerMiddleware + + broker = Broker(id="test-no-tracking") + assert not any(isinstance(mw, ContextTrackerMiddleware) for mw in broker.middlewares) + + +def test_tracking_enabled_registers_middleware(): + """settings.tracking.enabled=True auto-registers ContextTrackerMiddleware.""" + from moleculerpy.middleware.context_tracker import ContextTrackerMiddleware + from moleculerpy.settings import TrackingConfig + + settings = Settings(tracking=TrackingConfig(enabled=True, shutdown_timeout=2.5)) + broker = Broker(id="test-tracking", settings=settings) + + trackers = [mw for mw in broker.middlewares if isinstance(mw, ContextTrackerMiddleware)] + assert len(trackers) == 1 + # Seconds (2.5) -> milliseconds (2500) + assert trackers[0]._default_timeout == 2500 diff --git a/tests/unit/protocol_lifecycle_test.py b/tests/unit/protocol_lifecycle_test.py new file mode 100644 index 0000000..4bf2b2a --- /dev/null +++ b/tests/unit/protocol_lifecycle_test.py @@ -0,0 +1,296 @@ +"""System-level tests for the Protocol Correctness & Graceful Lifecycle sprint. + +These tests exercise the COMBINED behaviour of Tasks #1-#5: + +- Task #1: short broker hook aliases (starting/started/stopping) +- Task #2: PacketHeartbeat proto schema extended with seq/instanceID/memory/cpuSeq +- Task #3: Settings.tracking TrackingConfig +- Task #4: transit.send_disconnect_info + broker.stop() drain ordering +- Task #5: ContextTrackerMiddleware auto-registration via tracking.enabled + +The intent is to document the protocol contract end-to-end. Per-task unit +tests live next to their owning agent's changes; this file consolidates the +integration story. +""" + +from __future__ import annotations + +from unittest.mock import AsyncMock, Mock + +import pytest + +from moleculerpy.broker import ServiceBroker +from moleculerpy.middleware.base import Middleware +from moleculerpy.middleware.context_tracker import ContextTrackerMiddleware +from moleculerpy.settings import Settings, TrackingConfig +from moleculerpy.transit import Transit + +# --------------------------------------------------------------------------- +# Fixtures +# --------------------------------------------------------------------------- + + +@pytest.fixture +def mock_transit() -> AsyncMock: + """Mock Transit so brokers don't talk to real transports.""" + transit = AsyncMock(spec=Transit) + transit.connect = AsyncMock() + transit.disconnect = AsyncMock() + transit.send_disconnect_info = AsyncMock() + transit.ready = AsyncMock() + transit.send_node_info = AsyncMock() + transit.transporter = Mock(name="mock_transport") + return transit + + +def _make_broker( + mock_transit: AsyncMock, + *, + middlewares: list[Middleware] | None = None, + tracking: TrackingConfig | None = None, +) -> ServiceBroker: + settings = Settings(transporter="mock://localhost", tracking=tracking) + return ServiceBroker( + id="test-node", + settings=settings, + transit=mock_transit, + middlewares=middlewares or [], + ) + + +# --------------------------------------------------------------------------- +# Hook aliases (Task #1) +# --------------------------------------------------------------------------- + + +class ShortAliasMiddleware(Middleware): + """Node.js-style middleware using short alias names only.""" + + def __init__(self) -> None: + self.calls: list[str] = [] + + async def starting(self, broker): # type: ignore[override] + self.calls.append("starting") + + async def started(self, broker): # type: ignore[override] + self.calls.append("started") + + async def stopping(self, broker): # type: ignore[override] + self.calls.append("stopping") + + +class LongNameMiddleware(Middleware): + """Legacy MoleculerPy middleware using broker_* hook names.""" + + def __init__(self) -> None: + self.calls: list[str] = [] + + async def broker_starting(self, broker): # type: ignore[override] + self.calls.append("broker_starting") + + async def broker_started(self, broker): # type: ignore[override] + self.calls.append("broker_started") + + async def broker_stopping(self, broker): # type: ignore[override] + self.calls.append("broker_stopping") + + +class BothNamesMiddleware(Middleware): + """Middleware that defines BOTH broker_* and short alias variants.""" + + def __init__(self) -> None: + self.calls: list[str] = [] + + async def broker_started(self, broker): # type: ignore[override] + self.calls.append("broker_started") + + async def started(self, broker): # type: ignore[override] + self.calls.append("started") + + +@pytest.mark.asyncio +async def test_middleware_with_short_alias_invoked(mock_transit: AsyncMock) -> None: + mw = ShortAliasMiddleware() + broker = _make_broker(mock_transit, middlewares=[mw]) + await broker.start() + await broker.stop() + assert "starting" in mw.calls + assert "started" in mw.calls + assert "stopping" in mw.calls + + +@pytest.mark.asyncio +async def test_middleware_with_long_name_still_works(mock_transit: AsyncMock) -> None: + mw = LongNameMiddleware() + broker = _make_broker(mock_transit, middlewares=[mw]) + await broker.start() + await broker.stop() + assert mw.calls == [ + "broker_starting", + "broker_started", + "broker_stopping", + ] + + +@pytest.mark.asyncio +async def test_both_names_independent(mock_transit: AsyncMock) -> None: + mw = BothNamesMiddleware() + broker = _make_broker(mock_transit, middlewares=[mw]) + await broker.start() + await broker.stop() + # Both names get called — neither swallows the other + assert mw.calls.count("broker_started") == 1 + assert mw.calls.count("started") == 1 + + +# --------------------------------------------------------------------------- +# Heartbeat ProtoBuf roundtrip (Task #2) +# --------------------------------------------------------------------------- + + +def _protobuf_serializer(): + pytest.importorskip("google.protobuf") + from moleculerpy.serializers.protobuf import ProtoBufSerializer + + return ProtoBufSerializer() + + +def test_heartbeat_protobuf_roundtrip_with_seq() -> None: + serializer = _protobuf_serializer() + payload = { + "ver": "4", + "sender": "node-A", + "cpu": 12.5, + "seq": 42, + } + raw = serializer.serialize(payload, "HEARTBEAT") + decoded = serializer.deserialize(raw, "HEARTBEAT") + assert decoded.get("seq") == 42 + assert decoded.get("sender") == "node-A" + + +def test_heartbeat_protobuf_roundtrip_with_instanceid() -> None: + serializer = _protobuf_serializer() + payload = { + "ver": "4", + "sender": "node-A", + "cpu": 7.0, + "instanceID": "abc-123-instance", + "memory": 33.3, + "cpuSeq": 9, + } + raw = serializer.serialize(payload, "HEARTBEAT") + decoded = serializer.deserialize(raw, "HEARTBEAT") + assert decoded.get("instanceID") == "abc-123-instance" + assert decoded.get("cpuSeq") == 9 + # memory is a double in proto schema + assert decoded.get("memory") == pytest.approx(33.3, rel=1e-3) + + +def test_heartbeat_json_still_works() -> None: + from moleculerpy.serializers.json import JsonSerializer + + serializer = JsonSerializer() + payload = {"ver": "4", "sender": "n1", "cpu": 1.0, "seq": 7, "instanceID": "x"} + raw = serializer.serialize(payload, "HEARTBEAT") + decoded = serializer.deserialize(raw, "HEARTBEAT") + assert decoded == payload + + +# --------------------------------------------------------------------------- +# TrackingConfig (Task #3) +# --------------------------------------------------------------------------- + + +def test_tracking_config_import() -> None: + # Public re-export check — must be importable from moleculerpy.settings + import moleculerpy.settings as _settings + + cfg = _settings.TrackingConfig() + assert cfg.enabled is False + assert cfg.shutdown_timeout == 5.0 + + +def test_settings_default_tracking_disabled() -> None: + settings = Settings() + assert isinstance(settings.tracking, TrackingConfig) + assert settings.tracking.enabled is False + assert settings.tracking.shutdown_timeout == 5.0 + + +# --------------------------------------------------------------------------- +# Connection drain on stop (Task #4) +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_broker_stop_calls_send_disconnect_info(mock_transit: AsyncMock) -> None: + """send_disconnect_info must be called BEFORE transit.disconnect.""" + call_order: list[str] = [] + + async def _record_drain() -> None: + call_order.append("drain") + + async def _record_disconnect() -> None: + call_order.append("disconnect") + + mock_transit.send_disconnect_info.side_effect = _record_drain + mock_transit.disconnect.side_effect = _record_disconnect + + broker = _make_broker(mock_transit) + await broker.start() + await broker.stop() + + assert "drain" in call_order + assert "disconnect" in call_order + assert call_order.index("drain") < call_order.index("disconnect") + + +@pytest.mark.asyncio +async def test_send_disconnect_info_empty_services() -> None: + """Verify the drain INFO actually broadcasts services=[].""" + from moleculerpy.packet import Packet, Topic + + transit = Transit.__new__(Transit) # bypass __init__ + transit._was_connected = True # type: ignore[attr-defined] + transit.logger = Mock() + transit.publish = AsyncMock() # type: ignore[method-assign] + + fake_local_node = Mock() + fake_local_node.get_info.return_value = { + "sender": "node-A", + "services": [{"name": "math"}, {"name": "users"}], + "ver": "4", + } + transit.node_catalog = Mock() + transit.node_catalog.local_node = fake_local_node + + await Transit.send_disconnect_info(transit) + + transit.publish.assert_awaited_once() + pkt = transit.publish.await_args.args[0] + assert isinstance(pkt, Packet) + assert pkt.type == Topic.INFO + assert pkt.target is None + assert pkt.payload["services"] == [] + # Other fields preserved + assert pkt.payload["sender"] == "node-A" + + +# --------------------------------------------------------------------------- +# ContextTracker integration (Task #5) +# --------------------------------------------------------------------------- + + +def test_tracking_disabled_no_middleware(mock_transit: AsyncMock) -> None: + broker = _make_broker(mock_transit) # default tracking → disabled + assert not any(isinstance(mw, ContextTrackerMiddleware) for mw in broker.middlewares) + + +def test_tracking_enabled_adds_middleware(mock_transit: AsyncMock) -> None: + broker = _make_broker( + mock_transit, + tracking=TrackingConfig(enabled=True, shutdown_timeout=2.5), + ) + trackers = [mw for mw in broker.middlewares if isinstance(mw, ContextTrackerMiddleware)] + assert len(trackers) == 1 diff --git a/tests/unit/serializer_cbor_protobuf_test.py b/tests/unit/serializer_cbor_protobuf_test.py index 9aa7cb1..48971fd 100644 --- a/tests/unit/serializer_cbor_protobuf_test.py +++ b/tests/unit/serializer_cbor_protobuf_test.py @@ -197,6 +197,26 @@ def test_roundtrip_heartbeat(self, serializer: ProtoBufSerializer) -> None: result = serializer.deserialize(data, packet_type="HEARTBEAT") assert result["cpu"] == 75 + def test_roundtrip_heartbeat_extended_fields(self, serializer: ProtoBufSerializer) -> None: + # Verifies seq/instanceID/memory/cpuSeq survive proto roundtrip. + # These fields back the heartbeat-driven restart detection feature. + payload = { + "ver": "4", + "sender": "node-1", + "cpu": 50.5, + "seq": 42, + "instanceID": "abc-123-instance", + "memory": 1024.75, + "cpuSeq": 7, + } + data = serializer.serialize(payload, packet_type="HEARTBEAT") + result = serializer.deserialize(data, packet_type="HEARTBEAT") + assert result["seq"] == 42 + assert result["instanceID"] == "abc-123-instance" + assert result["memory"] == 1024.75 + assert result["cpuSeq"] == 7 + assert result["cpu"] == 50.5 + def test_roundtrip_ping_pong(self, serializer: ProtoBufSerializer) -> None: ping = {"ver": "4", "sender": "n1", "time": 1234567890, "id": "ping-1"} data = serializer.serialize(ping, packet_type="PING") diff --git a/tests/unit/settings_test.py b/tests/unit/settings_test.py index fa92960..7eab231 100644 --- a/tests/unit/settings_test.py +++ b/tests/unit/settings_test.py @@ -2,7 +2,7 @@ import pytest -from moleculerpy.settings import Settings, SettingsValidationError +from moleculerpy.settings import Settings, SettingsValidationError, TrackingConfig class TestSettings: @@ -331,3 +331,25 @@ def test_settings_valid_constants(self): assert "PLAIN" in Settings.VALID_LOG_FORMATS assert "JSON" in Settings.VALID_LOG_FORMATS + + +class TestTrackingConfig: + """Test TrackingConfig dataclass and Settings integration.""" + + def test_tracking_config_defaults(self): + s = Settings() + assert s.tracking.enabled is False + assert s.tracking.shutdown_timeout == 5.0 + + def test_tracking_config_custom(self): + s = Settings(tracking=TrackingConfig(enabled=True, shutdown_timeout=10.0)) + assert s.tracking.enabled is True + assert s.tracking.shutdown_timeout == 10.0 + + def test_tracking_config_validation(self): + with pytest.raises(SettingsValidationError) as exc_info: + Settings(tracking=TrackingConfig(shutdown_timeout=-1.0)) + assert "tracking.shutdown_timeout" in str(exc_info.value) + + with pytest.raises(SettingsValidationError): + Settings(tracking=TrackingConfig(shutdown_timeout=0)) diff --git a/tests/unit/transit_test.py b/tests/unit/transit_test.py index c015cf6..956121d 100644 --- a/tests/unit/transit_test.py +++ b/tests/unit/transit_test.py @@ -192,6 +192,51 @@ async def test_send_node_info(self, mock_dependencies, mock_transporter): assert packet.type == Topic.INFO assert packet.payload == {"id": "test-node", "services": []} + @pytest.mark.asyncio + async def test_send_disconnect_info(self, mock_dependencies, mock_transporter): + """send_disconnect_info broadcasts INFO with empty services list.""" + with patch("moleculerpy.transit.Transporter.get_by_name", return_value=mock_transporter): + transit = Transit(**mock_dependencies) + + # Not connected → no-op + transit._was_connected = False + await transit.send_disconnect_info() + mock_transporter.publish.assert_not_called() + + # Connected, with local node + transit._was_connected = True + mock_node = MagicMock() + mock_node.get_info.return_value = { + "id": "test-node", + "services": [{"name": "math"}], + "client": {"type": "python"}, + } + transit.node_catalog.local_node = mock_node + + await transit.send_disconnect_info() + + mock_transporter.publish.assert_called_once() + packet = mock_transporter.publish.call_args[0][0] + assert packet.type == Topic.INFO + assert packet.payload["services"] == [] + assert packet.payload["id"] == "test-node" + assert packet.payload["client"] == {"type": "python"} + + @pytest.mark.asyncio + async def test_send_disconnect_info_swallows_errors(self, mock_dependencies, mock_transporter): + """send_disconnect_info logs but doesn't raise on publish error.""" + with patch("moleculerpy.transit.Transporter.get_by_name", return_value=mock_transporter): + transit = Transit(**mock_dependencies) + transit._was_connected = True + mock_node = MagicMock() + mock_node.get_info.return_value = {"id": "n", "services": []} + transit.node_catalog.local_node = mock_node + mock_transporter.publish.side_effect = RuntimeError("boom") + + # Must not raise + await transit.send_disconnect_info() + transit.logger.warning.assert_called() + @pytest.mark.asyncio async def test_make_subscriptions(self, mock_dependencies, mock_transporter): """Test Transit _make_subscriptions method."""