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
218 changes: 115 additions & 103 deletions ventis/controller/global_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -534,120 +534,132 @@ def _poll_controllers(self):
if self.running:
self.process_supervisor.check_and_respawn()

for instance in self.instance_manager.list_instances():
# Polled in parallel, one instance's slow Redis/Postgres round-trip no longer
# gates every other instance's poll -- see ventis/OTLP_Exporter/DESIGN.md.
instances = self.instance_manager.list_instances()
if instances:
with ThreadPoolExecutor(max_workers=len(instances)) as executor:
list(executor.map(self._poll_one_instance, instances))

def _poll_one_instance(self, instance):
"""Poll and persist one instance's runtime/metrics/health data; never raises."""
try:
name = instance["agent_name"]
host = instance["host"]
port = instance["host_port"]
node_redis = self._get_node_redis_for(host)
try:
future_rows = pull_runtime_information(node_redis)
self._otel_db.write_waiting_rows(
future_rows, node_redis, self.config.get("project_id", 0)
)
except Exception as e:
logger.warning("Failed to poll instance %s: %s", instance, e)
return

# This is now legacy, keeping it for now, but will remove this later
send_runtime_information(
future_rows,
node_redis,
self.config.get("database", {}).get("url"),
)
except Exception as e:
logger.warning(
"Failed to write runtime information for instance %s (%s:%s) "
"(non-fatal): %s",
name,
host,
port,
e,
try:
future_rows = pull_runtime_information(node_redis)
self._otel_db.write_waiting_rows(
future_rows, node_redis, self.config.get("project_id", 0)
)
# This is now legacy, keeping it for now, but will remove this later
send_runtime_information(
future_rows,
node_redis,
self.config.get("database", {}).get("url"),
)
except Exception as e:
logger.warning(
"Failed to write runtime information for instance %s (%s:%s) "
"(non-fatal): %s",
name,
host,
port,
e,
)
agent_host = self._agent_host_key(host)
status_key = f"controller:{agent_host}:{port}:status"
metrics_key = f"controller:{agent_host}:{port}:metrics"

# Getting metrics from local controllers
# See LocalController._execute_locally
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
)
agent_host = self._agent_host_key(host)
status_key = f"controller:{agent_host}:{port}:status"
metrics_key = f"controller:{agent_host}:{port}:metrics"
throughput = requests_served / elapsed if elapsed > 0 else 0.0
self._last_metrics_poll_time[(host, port)] = now

# Getting metrics from local controllers
# See LocalController._execute_locally
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
try:
send_agent_information(
[
{
**instance,
**metrics,
"requests_served": requests_served,
"throughput": throughput,
}
],
self.config.get("database", {}).get("url"),
)
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,
)
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,
)

try:
status = node_redis.get(status_key) or "unknown"
prev = self._last_status.get((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
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:
# 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,
)
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)
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