diff --git a/ventis/controller/global_controller.py b/ventis/controller/global_controller.py index ba53881..6bdbd1a 100644 --- a/ventis/controller/global_controller.py +++ b/ventis/controller/global_controller.py @@ -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 #