Skip to content
Merged
22 changes: 22 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.20] - 2026-04-06

### Added
Expand Down
2 changes: 1 addition & 1 deletion moleculerpy/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@
try:
__version__ = version("moleculerpy")
except PackageNotFoundError:
__version__ = "0.14.20"
__version__ = "0.14.21"

__all__ = [ # noqa: RUF022
# Core
Expand Down
67 changes: 44 additions & 23 deletions moleculerpy/cacher/redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -124,9 +124,12 @@ def __init__(
# Default TTL in seconds
self.default_ttl: int | None = self.opts.get("ttl")

# Serializer (initialized in init())
# Serializer
self.serializer = JSONSerializer()

# Logger fallback (overridden in init() with broker logger)
self.logger = logger

# Ping interval (optional periodic health check)
self.ping_interval = self.opts.get("ping_interval")
self._ping_task: asyncio.Task[None] | None = None
Expand All @@ -139,16 +142,16 @@ def init(self, broker: ServiceBroker) -> None:
"""
super().init(broker)

# Redis is NOT connected until start() → connect() is called.
# BaseCacher.init() sets connected=True (fine for MemoryCacher), but
# for network cachers we override back to False until actual connection.
self.connected = False

# Create logger (use broker's logger factory)
if hasattr(broker, "_create_logger"):
self.logger = broker._create_logger("REDIS-CACHER")
else:
self.logger = logger

# Add namespace to prefix if configured
if broker.namespace and broker.namespace != "":
self.prefix = f"MOL-{broker.namespace}-"

# Namespace prefix handled by BaseCacher.init() — no duplicate logic here.
self.logger.info(f"Initializing Redis cacher with prefix '{self.prefix}'")

async def connect(self) -> None:
Expand Down Expand Up @@ -228,24 +231,25 @@ async def disconnect(self) -> None:
self.connected = False

async def _ping_loop(self) -> None:
"""Periodic ping task for connection monitoring."""
"""Periodic ping task for connection monitoring.

Pings first, then sleeps — detects immediate connection drop.
"""
try:
interval = self.ping_interval if isinstance(self.ping_interval, (int, float)) else 10.0
while self.connected and self.client:
interval = (
self.ping_interval if isinstance(self.ping_interval, (int, float)) else 0.0
)
await asyncio.sleep(interval)
try:
await self._await_maybe(self.client.ping())
except Exception as e:
self.logger.warning(f"Ping failed: {e}")
self.logger.error(f"Redis ping failed — connection may be broken: {e}")
self.connected = False
# Broadcast error event
if self.broker:
await self.broker.broadcast_local(
"$cacher.error",
{"error": str(e), "module": "cacher", "type": "PING_FAILED"},
)
break
await asyncio.sleep(interval)
except asyncio.CancelledError:
pass

Expand Down Expand Up @@ -329,21 +333,25 @@ async def set(

prefixed_key = self._get_prefixed_key(key)

# Use provided TTL or default
if ttl is None:
ttl = self.default_ttl
# Use provided TTL or default; validate
effective_ttl = ttl if ttl is not None else self.default_ttl
if effective_ttl is not None and effective_ttl <= 0:
self.logger.warning(
f"Invalid TTL {effective_ttl} for key {key}, storing without expiry"
)
effective_ttl = None

try:
# Serialize data
serialized = self.serializer.serialize(data)

# SET with optional TTL (atomic operation)
if ttl:
await self.client.set(prefixed_key, serialized, ex=ttl)
if effective_ttl is not None:
await self.client.set(prefixed_key, serialized, ex=effective_ttl)
else:
await self.client.set(prefixed_key, serialized)

self.logger.debug(f"Cache SET: {key} (ttl={ttl}s)")
self.logger.debug(f"Cache SET: {key} (ttl={effective_ttl}s)")

except Exception as e:
self.logger.error(f"Redis SET error for {key}: {e}")
Expand Down Expand Up @@ -517,9 +525,22 @@ async def get_cache_keys(self) -> list[dict[str, str]]:
self.logger.error(f"Redis GET_CACHE_KEYS error: {e}")
return []

# Lock methods (TODO: Implement Redlock for distributed locking)
# For now, inherit from BaseCacher which provides in-memory lock fallback
async def start(self) -> None:
"""Start cacher — connect to Redis.

Called by broker during broker.start(). Matches Node.js init() behavior
where connection is established during broker lifecycle.
"""
await self.connect()

async def stop(self) -> None:
"""Stop cacher — disconnect from Redis.

Called by broker during broker.stop().
"""
await super().stop() # Clears in-memory locks
await self.disconnect()

async def close(self) -> None:
"""Close Redis connection gracefully."""
"""Close Redis connection gracefully (alias for disconnect)."""
await self.disconnect()
7 changes: 7 additions & 0 deletions moleculerpy/discoverer.py
Original file line number Diff line number Diff line change
Expand Up @@ -203,6 +203,13 @@ def check_remote_nodes(self) -> None:
)
node_catalog.disconnect_node(node_id, unexpected=True)

# Evict stale entries from _discover_pending (prevents unbounded growth)
expired = [
k for k, ts in self._discover_pending.items() if now - ts > self._DISCOVER_COOLDOWN
]
for k in expired:
self._discover_pending.pop(k, None)

def check_offline_nodes(self) -> None:
"""Check offline nodes. Remove which are older than 10 minutes.

Expand Down
1 change: 1 addition & 0 deletions moleculerpy/node.py
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,7 @@ def process_node_info(self, node_id: str, payload: dict[str, Any]) -> None:
node.client = payload.get("client")
node.metadata = payload.get("metadata", {})
node.seq = payload.get("seq", 0)
node.instanceID = payload.get("instanceID", node.instanceID)
node.port = payload.get("port", 0)

# Phase 5.1: Update extended metrics
Expand Down
31 changes: 30 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 @@ -577,6 +582,8 @@ async def _handle_heartbeat(self, packet: Packet) -> None:
"""
if not packet.sender or packet.sender == self.node_id:
return # Ignore own heartbeats (Node.js: sender === this.broker.nodeID)
if not isinstance(packet.payload, dict):
return # Malformed heartbeat

node = self.node_catalog.get_node(packet.sender)
if node is None:
Expand All @@ -587,6 +594,28 @@ 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 int(getattr(node, "seq", 0)) != int(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)
# Skip if node has no instanceID yet (first registration, no INFO received)
payload_iid = packet.payload.get("instanceID")
node_iid = getattr(node, "instanceID", None)
if payload_iid is not None and node_iid 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
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "hatchling.build"

[project]
name = "moleculerpy"
version = "0.14.20"
version = "0.14.21"
description = "Fast, modern microservices framework for Python - Port of Moleculer.js"
authors = [
{ name = "Eli Rum", email = "explosivebit@gmail.com" }
Expand Down
Loading
Loading