diff --git a/moleculerpy/transit.py b/moleculerpy/transit.py index f475034..3ed77a6 100644 --- a/moleculerpy/transit.py +++ b/moleculerpy/transit.py @@ -519,11 +519,16 @@ async def beat(self) -> None: local_node.hostname = static["hostname"] local_node.ipList = static["ip_list"] - heartbeat_data = { + heartbeat_data: dict[str, Any] = { "cpu": metrics["cpu"], "cpuSeq": metrics["cpuSeq"], "memory": metrics["memory"], # Python extension } + # Include seq and instanceID so remote nodes can detect service changes + # and restarts via heartbeat (Node.js checks these in heartbeatReceived). + if local_node: + heartbeat_data["seq"] = local_node.seq + heartbeat_data["instanceID"] = local_node.instanceID await self.publish(Packet(Topic.HEARTBEAT, None, heartbeat_data)) async def send_node_info(self) -> None: @@ -587,6 +592,27 @@ async def _handle_heartbeat(self, packet: Packet) -> None: await self._request_discovery(packet.sender, "offline") return + # Check seq mismatch — services changed on remote node (Node.js parity) + payload_seq = packet.payload.get("seq") + if payload_seq is not None and getattr(node, "seq", 0) != payload_seq: + self.logger.debug( + f"Service seq changed on '{packet.sender}' " + f"({getattr(node, 'seq', 0)} → {payload_seq}), requesting INFO" + ) + await self._request_discovery(packet.sender, "seq-changed") + return + + # Check instanceID mismatch — node restarted (Node.js parity) + payload_iid = packet.payload.get("instanceID") + node_iid = getattr(node, "instanceID", None) or "" + if payload_iid is not None and not str(node_iid).startswith(str(payload_iid)): + self.logger.debug( + f"instanceID changed on '{packet.sender}' " + f"({node_iid} → {payload_iid}), requesting INFO" + ) + await self._request_discovery(packet.sender, "instance-restarted") + return + # Update Moleculer.js compatible metrics node.cpu = packet.payload.get("cpu", 0) node.cpuSeq = packet.payload.get("cpuSeq", 0) diff --git a/tests/unit/transit_test.py b/tests/unit/transit_test.py index 9b551cd..c015cf6 100644 --- a/tests/unit/transit_test.py +++ b/tests/unit/transit_test.py @@ -1678,6 +1678,74 @@ async def test_heartbeat_no_sender_ignored(self, mock_dependencies, mock_transpo transit.discover_node.assert_not_awaited() transit.node_catalog.get_node.assert_not_called() + @pytest.mark.asyncio + async def test_heartbeat_seq_mismatch_triggers_discover( + self, mock_dependencies, mock_transporter + ): + """Heartbeat with different seq triggers re-discovery (services changed).""" + with patch("moleculerpy.transit.Transporter.get_by_name", return_value=mock_transporter): + transit = Transit(**mock_dependencies) + mock_node = MagicMock() + mock_node.available = True + mock_node.seq = 5 + mock_node.instanceID = "abc-123" + transit.node_catalog = MagicMock() + transit.node_catalog.get_node.return_value = mock_node + transit._request_discovery = AsyncMock() + + packet = Packet(Topic.HEARTBEAT, "remote", {"cpu": 50, "seq": 10}) + packet.sender = "remote" + await transit._handle_heartbeat(packet) + + transit._request_discovery.assert_awaited_once_with("remote", "seq-changed") + + @pytest.mark.asyncio + async def test_heartbeat_instanceid_mismatch_triggers_discover( + self, mock_dependencies, mock_transporter + ): + """Heartbeat with different instanceID triggers re-discovery (node restarted).""" + with patch("moleculerpy.transit.Transporter.get_by_name", return_value=mock_transporter): + transit = Transit(**mock_dependencies) + mock_node = MagicMock() + mock_node.available = True + mock_node.seq = 5 + mock_node.instanceID = "abc-123" + transit.node_catalog = MagicMock() + transit.node_catalog.get_node.return_value = mock_node + transit._request_discovery = AsyncMock() + + packet = Packet( + Topic.HEARTBEAT, "remote", {"cpu": 50, "seq": 5, "instanceID": "xyz-999"} + ) + packet.sender = "remote" + await transit._handle_heartbeat(packet) + + transit._request_discovery.assert_awaited_once_with("remote", "instance-restarted") + + @pytest.mark.asyncio + async def test_heartbeat_same_seq_instanceid_updates_metrics( + self, mock_dependencies, mock_transporter + ): + """Heartbeat with matching seq/instanceID just updates metrics (no discovery).""" + with patch("moleculerpy.transit.Transporter.get_by_name", return_value=mock_transporter): + transit = Transit(**mock_dependencies) + mock_node = MagicMock() + mock_node.available = True + mock_node.seq = 5 + mock_node.instanceID = "abc-123" + transit.node_catalog = MagicMock() + transit.node_catalog.get_node.return_value = mock_node + transit._request_discovery = AsyncMock() + + packet = Packet( + Topic.HEARTBEAT, "remote", {"cpu": 80, "seq": 5, "instanceID": "abc-123"} + ) + packet.sender = "remote" + await transit._handle_heartbeat(packet) + + transit._request_discovery.assert_not_awaited() + assert mock_node.cpu == 80 + class TestRequestDiscovery: """Tests for _request_discovery rate-limiting and _handle_discover targeted reply."""