From a53afe5aa56eee3e52053b66d06286f85e32765f Mon Sep 17 00:00:00 2001 From: Nick Huo Date: Tue, 1 Sep 2026 22:02:31 -0700 Subject: [PATCH] [CAN-283] Skip telemetry writes when no database is configured With no `database` in the controller config, every metrics poll built an engine from the string "None", so SQLAlchemy raised "Could not parse SQLAlchemy URL from given URL string" once per instance every 5 seconds. Telemetry has nowhere to go without a database, so the writes -- and the metrics that only exist to feed them -- are now skipped, and the fact is stated once at startup. --- .../test_global_controller_telemetry_skip.py | 131 +++++++++++++++ ventis/controller/global_controller.py | 149 ++++++++++-------- ventis/controller/utils/telemetry_logging.py | 9 +- 3 files changed, 222 insertions(+), 67 deletions(-) create mode 100644 tests/test_global_controller_telemetry_skip.py diff --git a/tests/test_global_controller_telemetry_skip.py b/tests/test_global_controller_telemetry_skip.py new file mode 100644 index 0000000..2c0e711 --- /dev/null +++ b/tests/test_global_controller_telemetry_skip.py @@ -0,0 +1,131 @@ +"""CAN-283: with no `database` in the config, every metrics poll used to log + + Failed to write agent information for instance (localhost:800N) (non-fatal): + Could not parse SQLAlchemy URL from given URL string + +once per instance, every poll_interval (5s by default), drowning out the errors an +operator actually needs to read. Telemetry has nowhere to go without a database, so +the writes -- and the metrics that feed them -- are skipped outright. +""" + +import os +import sys +import unittest +from unittest.mock import MagicMock, patch + +sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))) + +from ventis.controller.global_controller import GlobalController +from ventis.controller.utils.telemetry_logging import resolve_database_url + + +class EnvIsolatedTestCase(unittest.TestCase): + def setUp(self): + self._saved = os.environ.pop("VENTIS_DATABASE_URL", None) + + def tearDown(self): + os.environ.pop("VENTIS_DATABASE_URL", None) + if self._saved is not None: + os.environ["VENTIS_DATABASE_URL"] = self._saved + + +class ResolveDatabaseUrlTests(EnvIsolatedTestCase): + def test_nothing_configured_resolves_to_none(self): + # config.get("database", {}).get("url") yields None when the key is absent -- + # str(None) == "None" is what SQLAlchemy used to choke on. + for value in (None, "", " "): + self.assertIsNone(resolve_database_url(value), repr(value)) + + def test_env_var_wins_over_the_config_value(self): + os.environ["VENTIS_DATABASE_URL"] = "sqlite:///env.db" + self.assertEqual( + resolve_database_url("postgresql://cfg-host/db"), "sqlite:///env.db" + ) + + def test_empty_env_var_falls_back_to_the_config_value(self): + os.environ["VENTIS_DATABASE_URL"] = "" + self.assertEqual(resolve_database_url("sqlite:///cfg.db"), "sqlite:///cfg.db") + + +class PollTelemetryTests(EnvIsolatedTestCase): + """_poll_controllers must not touch telemetry when there is no database URL.""" + + def _controller(self, config): + controller = GlobalController.__new__(GlobalController) + controller.config = config + controller.poll_interval = 5 + controller._last_status = {} + controller._last_metrics_poll_time = {} + controller._on_controller_healthy = MagicMock() + controller.instance_manager = MagicMock() + controller.instance_manager.list_instances.return_value = [ + { + "agent_name": "AgentA", + "agent_id": "local:AgentA:0", + "host": "localhost", + "host_port": 8000, + } + ] + self.node_redis = MagicMock() + self.node_redis.hgetall.return_value = {"requests_served": "3"} + self.node_redis.get.return_value = "healthy" + controller._get_node_redis_for = MagicMock(return_value=self.node_redis) + return controller + + def _poll(self, controller): + with patch( + "ventis.controller.global_controller.send_runtime_information" + ) as send_runtime, patch( + "ventis.controller.global_controller.send_agent_information" + ) as send_agent, patch( + "ventis.controller.global_controller.pull_runtime_information" + ) as pull_runtime, self.assertLogs( + "ventis.controller.global_controller", level="INFO" + ) as logs: + controller._poll_controllers() + return send_runtime, send_agent, pull_runtime, logs.output + + def _assert_telemetry_skipped(self, config): + controller = self._controller(config) + + send_runtime, send_agent, pull_runtime, output = self._poll(controller) + + send_runtime.assert_not_called() + send_agent.assert_not_called() + pull_runtime.assert_not_called() + self.assertEqual([line for line in output if "WARNING" in line], []) + # Nothing was persisted, so the counters stay where they are. + self.node_redis.hset_multiple.assert_not_called() + + def test_a_missing_database_key_skips_telemetry(self): + self._assert_telemetry_skipped({"agents": []}) + + def test_an_empty_database_block_is_treated_as_no_database(self): + # `database:` with nothing under it parses as None, not {}. + self._assert_telemetry_skipped({"agents": [], "database": None}) + + def test_a_configured_database_still_gets_both_writes(self): + controller = self._controller( + {"agents": [], "database": {"url": "postgresql://user@host/db"}} + ) + + send_runtime, send_agent, _, output = self._poll(controller) + + self.assertEqual(send_runtime.call_args.args[2], "postgresql://user@host/db") + self.assertEqual(send_agent.call_args.args[1], "postgresql://user@host/db") + self.assertEqual([line for line in output if "WARNING" in line], []) + # Counters are cleared only once the row has actually been persisted. + self.node_redis.hset_multiple.assert_called_once() + + def test_the_env_var_alone_is_enough_to_enable_writes(self): + os.environ["VENTIS_DATABASE_URL"] = "sqlite:///env.db" + controller = self._controller({"agents": []}) + + send_runtime, send_agent, _, _ = self._poll(controller) + + send_runtime.assert_called_once() + send_agent.assert_called_once() + + +if __name__ == "__main__": + unittest.main() diff --git a/ventis/controller/global_controller.py b/ventis/controller/global_controller.py index 1e24f10..21fa260 100644 --- a/ventis/controller/global_controller.py +++ b/ventis/controller/global_controller.py @@ -20,6 +20,7 @@ from ventis.controller.utils.telemetry_logging import ( assign_project_id, pull_runtime_information, + resolve_database_url, send_runtime_information, send_agent_information, ) @@ -84,6 +85,12 @@ def __init__(self, config_path): self._lc_stubs = {} # endpoint -> gRPC stub self.instance_manager = InstanceManager(self) assign_project_id(self.config.get("project_id",0)) + if self._database_url() is None: + logger.info( + "No database configured; telemetry writes are disabled. " + "Set database.url in %s to record runtime and agent information.", + config_path, + ) # Clean up any stale containers from previous runs self._cleanup_stale_containers() @@ -440,83 +447,32 @@ def run(self): except KeyboardInterrupt: self.stop() + def _database_url(self): + """The configured database URL, or None when there is no database to write to. + + Read per call, not cached, so reload_config() can point us at a new database. + `database:` with nothing under it parses as None, hence the `or {}`. + """ + return resolve_database_url((self.config.get("database") or {}).get("url")) + def _poll_controllers(self): """ Check the health of each registered controller replica via its node's Redis. Also retrieves the request calls made in each instance. """ + # Without a database the telemetry writes -- and the metrics that feed them -- + # have nowhere to go, so they are skipped instead of failing on every poll. + database_url = self._database_url() + for instance in self.instance_manager.list_instances(): name = instance["agent_name"] host = instance["host"] port = instance["host_port"] node_redis = self._get_node_redis_for(host) - try: - send_runtime_information( - pull_runtime_information(node_redis), - 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" + if database_url: + self._write_telemetry(instance, node_redis, database_url) - # 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 - ) - 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_key = f"controller:{self._agent_host_key(host)}:{port}:status" try: status = node_redis.get(status_key) or "unknown" @@ -554,6 +510,67 @@ def _poll_controllers(self): e, ) + def _write_telemetry(self, instance, node_redis, database_url): + """Persist runtime and per-instance metrics for one instance. Never fatal.""" + name = instance["agent_name"] + host = instance["host"] + port = instance["host_port"] + + try: + send_runtime_information( + pull_runtime_information(node_redis), + node_redis, + database_url, + ) + except Exception as e: + logger.warning( + "Failed to write runtime information for instance %s (%s:%s) " + "(non-fatal): %s", + name, + host, + port, + e, + ) + + # Getting metrics from local controllers + # See LocalController._execute_locally + metrics_key = f"controller:{self._agent_host_key(host)}:{port}:metrics" + try: + metrics = node_redis.hgetall(metrics_key) + if not metrics: + return + 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 + ) + self._last_metrics_poll_time[(host, port)] = now + send_agent_information( + [ + { + **instance, + **metrics, + "requests_served": requests_served, + "throughput": requests_served / elapsed if elapsed > 0 else 0.0, + } + ], + database_url, + ) + except Exception as e: + logger.warning( + "Failed to record metrics 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}, + ) + # ------------------------------------------------------------------ # # Extensibility hooks — override in subclasses # # ------------------------------------------------------------------ # diff --git a/ventis/controller/utils/telemetry_logging.py b/ventis/controller/utils/telemetry_logging.py index 7d911e2..0100b44 100644 --- a/ventis/controller/utils/telemetry_logging.py +++ b/ventis/controller/utils/telemetry_logging.py @@ -87,10 +87,17 @@ def assign_project_id(project_id) -> None: global _project_id _project_id = project_id + +def resolve_database_url(database_url): + """The env var wins over the config value, as it always has; empty means no database.""" + url = os.environ.get("VENTIS_DATABASE_URL") or str(database_url or "") + return url.strip() or None + + def _get_engine(database_url): global _engine if _engine is None: - url = os.environ.get("VENTIS_DATABASE_URL", str(database_url)) + url = resolve_database_url(database_url) or "" if url.startswith("postgresql://"): url = "postgresql+psycopg://" + url[len("postgresql://"):] _engine = create_engine(url)