Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 27 additions & 1 deletion moleculerpy/transit.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down
68 changes: 68 additions & 0 deletions tests/unit/transit_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
Loading