Skip to content
Open
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
140 changes: 79 additions & 61 deletions ventis/controller/global_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -433,70 +433,88 @@ def _poll_controllers(self):

# Getting metrics from local controllers
# See LocalController._execute_locally
metrics = node_redis.hgetall(metrics_key)
if metrics:
now = time.time()
requests_served = int(float(metrics.get("requests_served") or 0))
elapsed = now - self._last_metrics_poll_time.get(
(host, port), now - self.poll_interval
)
throughput = requests_served / elapsed if elapsed > 0 else 0.0
self._last_metrics_poll_time[(host, port)] = now

try:
send_agent_information(
[
{
**instance,
**metrics,
"requests_served": requests_served,
"throughput": throughput,
}
],
self.config.get("database", {}).get("url"),
)
except Exception as e:
logger.warning(
"Failed to write agent information for instance %s (%s:%s) "
"(non-fatal): %s",
name,
host,
port,
e,
)
else:
# Only clear the accumulated counters once they've actually been persisted
node_redis.hset_multiple(
metrics_key,
{"full_failures": 0, "error_count": 0, "requests_served": 0},
try:
metrics = node_redis.hgetall(metrics_key)
if metrics:
now = time.time()
requests_served = int(float(metrics.get("requests_served") or 0))
elapsed = now - self._last_metrics_poll_time.get(
(host, port), now - self.poll_interval
)
throughput = requests_served / elapsed if elapsed > 0 else 0.0
self._last_metrics_poll_time[(host, port)] = now

try:
send_agent_information(
[
{
**instance,
**metrics,
"requests_served": requests_served,
"throughput": throughput,
}
],
self.config.get("database", {}).get("url"),
)
except Exception as e:
logger.warning(
"Failed to write agent information for instance %s (%s:%s) "
"(non-fatal): %s",
name,
host,
port,
e,
)
else:
# Only clear the accumulated counters once they've actually been persisted
node_redis.hset_multiple(
metrics_key,
{"full_failures": 0, "error_count": 0, "requests_served": 0},
)
except Exception as e:
logger.warning(
"Failed to poll metrics for instance %s (%s:%s): %s",
name,
host,
port,
e,
)

status = node_redis.get(status_key) or "unknown"
prev = self._last_status.get((host, port))

if status != prev:
if status == "healthy":
logger.info(
"Controller %s (%s:%s) is now healthy.", name, host, port
)
self._on_controller_healthy(name, host, port)
else:
logger.warning(
"Controller %s (%s:%s) status changed: %s -> %s",
name,
host,
port,
prev or "(none)",
status,
)
self._on_controller_unhealthy(name, host, port)
self._last_status[(host, port)] = status
else:
# No change — healthy stays quiet, unhealthy stays quiet too
if status == "healthy":
self._on_controller_healthy(name, host, port)
try:
status = node_redis.get(status_key) or "unknown"
prev = self._last_status.get((host, port))

if status != prev:
if status == "healthy":
logger.info(
"Controller %s (%s:%s) is now healthy.", name, host, port
)
self._on_controller_healthy(name, host, port)
else:
logger.warning(
"Controller %s (%s:%s) status changed: %s -> %s",
name,
host,
port,
prev or "(none)",
status,
)
self._on_controller_unhealthy(name, host, port)
self._last_status[(host, port)] = status
else:
self._on_controller_unhealthy(name, host, port)
# No change — healthy stays quiet, unhealthy stays quiet too
if status == "healthy":
self._on_controller_healthy(name, host, port)
else:
self._on_controller_unhealthy(name, host, port)
except Exception as e:
logger.warning(
"Failed to poll status for instance %s (%s:%s): %s",
name,
host,
port,
e,
)

# ------------------------------------------------------------------ #
# Extensibility hooks — override in subclasses #
Expand Down
Loading