diff --git a/coco/config.py b/coco/config.py index c50a7728..1b7e57a3 100644 --- a/coco/config.py +++ b/coco/config.py @@ -25,9 +25,6 @@ host: localhost port: 12055 - # Port for prometheus metrics - metrics_port: 12056 - # Port the redis server is listening on. redis_port: 6379 @@ -124,7 +121,6 @@ _config_skeleton = { "host": RequiredValue, "port": DEFAULT_PORT, - "metrics_port": 9090, "redis_port": 6379, "log_level": "INFO", "endpoint_dir": RequiredValue, @@ -144,6 +140,9 @@ "comet_broker": {"enabled": True}, } +# List of config keys to warn about if they're present +_warn_if_present = {"metrics_port"} + def load_config( path: str | os.PathLike | None = None, cli: bool = False, testing: bool = False @@ -294,6 +293,8 @@ def _validate_and_resolve(config: dict) -> None: for key, value in config.items(): if value is RequiredValue: missing_values.append(key) + elif key in _warn_if_present: + logger.warning(f'Unused config key "{key!r}" will be ignored.') if missing_values: raise click.ClickException( diff --git a/coco/core.py b/coco/core.py index 6b79906d..43b8347a 100644 --- a/coco/core.py +++ b/coco/core.py @@ -12,7 +12,7 @@ import os import socket import time -from multiprocessing import set_start_method +from multiprocessing import Pipe, set_start_method from pathlib import Path import click @@ -20,7 +20,7 @@ from comet import CometError, Manager from sanic import Sanic, response -from . import config, slack, wait, worker +from . import config, metrics, slack, wait, worker from .endpoint import ( Endpoint, LocalEndpoint, @@ -39,6 +39,17 @@ logger = logging.getLogger(__name__) +# These are all the local endpoints +all_local_endpoints = { + "blocklist", + "update-blocklist", + "saved-states", + "reset-state", + "save-state", + "load-state", + "wait", +} + # This should be a no-op on Linux but is required on MacOS for coco to run try: set_start_method("fork", force=True) @@ -108,16 +119,6 @@ def __init__( all_endpoints = {endpoint["name"] for endpoint in self.config["endpoints"]} - # These are all the local endpoints - all_local_endpoints = { - "blocklist", - "update-blocklist", - "saved-states", - "reset-state", - "save-state", - "load-state", - "wait", - } all_endpoints |= all_local_endpoints for endpoint in self.config["endpoints"]: @@ -126,6 +127,9 @@ def __init__( ) self.config["endpoints"] = endpoints + # Pipe to communicate with the metrics aggregator + self.mpipe = None + # Configure the forwarder try: timeout = str2total_seconds(self.config["timeout"]) @@ -133,8 +137,10 @@ def __init__( raise click.ClickException( f"Failed parsing value 'timeout' ({self.config['timeout']}): {e}" ) from e + self.forwarder = RequestForwarder( self.blocklist_path, + self.config["redis_port"], timeout, debug_connections=self.config["debug_connections"], ) @@ -145,7 +151,7 @@ def __init__( self._config_slack_loggers() self._load_endpoints() - self._local_endpoints(all_local_endpoints) + self._local_endpoints() self._check_endpoint_links() try: @@ -166,6 +172,14 @@ def __init__( self.redis_sync = redis.Redis(port=int(self.config["redis_port"])) self.redis_sync.lrem("queue", 0, "coco_shutdown") + # Remove old metric keys + self.redis_sync.delete( + "coco_calls", + "dropped_requests", + "external_response_time", + "queue_wait_time", + ) + # Load queue update script into redis cache self.queue_sha = self.redis_sync.script_load( """ if redis.call('llen', KEYS[1]) >= tonumber(ARGV[1]) then @@ -192,22 +206,17 @@ def __init__( sock = None # This blocks until cocod terminates - self._start_server(sock) + self._start_server(all_endpoints, sock) # Coco is done, delete the redis async connection self.redis_async = None def _call_endpoints_on_start(self): - """This is a short-lived Sanic worker. - - It takes care of initialising redis for the endpoints - and handles "call_on_start" endpoints.""" + """A short-lived Sanic worker to handle "call_on_start" endpoints.""" logger.debug("init-endpoints worker start-up") for endpoint in self.endpoints.values(): - # Initialise request counter - self.redis_sync.incr(f"dropped_counter_{endpoint.name}", amount=0) if endpoint.call_on_start: logger.debug(f"Calling endpoint on start: /{endpoint.name}") name = f"{os.getpid()}-{time.time()}" @@ -237,11 +246,13 @@ def _call_endpoints_on_start(self): logger.debug("init-endpoints worker finished (exiting)") - def _start_server(self, sock: socket.socket | None = None): + def _start_server(self, all_endpoints: set, sock: socket.socket | None = None): """Start a sanic server. Parameters ---------- + all_endpoints : set + A set of all the endpoint names sock : socket.socket or None A bound socket for Sanic to listen on, if in --testing mode, or None, in production mode. @@ -252,6 +263,25 @@ def _start_server(self, sock: socket.socket | None = None): def start_workers(app): """Start the non-Sanic worker processes.""" + nonlocal all_endpoints + + # Create a Pipe for communication with the metrics aggregator + self.mpipe, aggregator_pipe_end = Pipe() + + # Start the metrics aggregator + app.manager.manage( + "metrics-aggregator", + metrics.aggregator, + { + "app": app, + "all_endpoints": all_endpoints, + "pipe": aggregator_pipe_end, + "port": int(self.config["redis_port"]), + }, + transient=False, + restartable=True, + auto_start=True, + ) # Start the qworker (the cocod back-end) app.manager.manage( @@ -260,9 +290,7 @@ def start_workers(app): { "app": app, "endpoints": self.endpoints, - "forwarder": self.forwarder, "coco_port": self.config["port"], - "metrics_port": self.config["metrics_port"], "redis_port": int(self.config["redis_port"]), "frontend_timeout": self.frontend_timeout, }, @@ -293,7 +321,7 @@ def signal_coco_shutdown(app): except RuntimeError as e: logger.error(f"queueing coco_shutdown in redis failed: {e}") - # Get Sanic to start/stop the qworker process when it starts/terminates + # Listeners to start (and stop) the non-Sanic worker processes self.sanic_app.register_listener(start_workers, "main_process_ready") self.sanic_app.register_listener(signal_coco_shutdown, "before_server_stop") @@ -324,6 +352,7 @@ async def stop_slack_log(*_): # non-endpoint routes. These take precedence over an identically named # endpoint self.sanic_app.add_route(self._get_config, "/config", methods=["GET"]) + self.sanic_app.add_route(self._send_metrics, "/metrics", methods=["GET"]) self.sanic_app.add_route(self._return_qlen, "/qlen", methods=["GET"]) self.sanic_app.add_route( @@ -482,7 +511,7 @@ def _load_endpoints(self): ) self.forwarder.add_endpoint(name, self.endpoints[name]) - def _local_endpoints(self, all_local_endpoints): + def _local_endpoints(self): # Register any local endpoints endpoints = { @@ -510,7 +539,29 @@ async def _get_config(self, _): """ return response.json(self.config) + async def _send_metrics(self, _): + """Sanic handler for the /metrics route. + + Fetches the rendered metrics from the metrics + aggregator and forwards them on to the client. + """ + if not self.mpipe: + return response.text("") + + # Request metrics from the aggregator. Anything sent + # to the aggregator over the pipe triggers it to respond + # with the rendered metrics. + self.mpipe.send(True) + + # The aggregregator responds with a UTF-8-encoded body and a content type + body, content_type = self.mpipe.recv() + return response.raw(body, content_type=content_type) + async def _return_qlen(self, _): + """Sanic handler for the /qlen route. + + Returns the length of the queue (as plain text). + """ from sanic import text return text(str(self.redis_sync.llen("queue"))) @@ -582,8 +633,8 @@ async def external_endpoint(self, request, endpoint): ) if full: - # Increment dropped request counter - await ra_cli.incr(f"dropped_counter_{endpoint}") + # Append to the dropped_request list + await ra_cli.rpush("dropped_requests", endpoint) return response.json( {"reply": "Coco queue is full.", "status": 503}, status=503 ) diff --git a/coco/metric.py b/coco/metric.py deleted file mode 100644 index 5d074163..00000000 --- a/coco/metric.py +++ /dev/null @@ -1,93 +0,0 @@ -""" -coco metric module. - -Helper functions for prometheus metric exporting. -""" - -import logging -import threading -from typing import ClassVar - -import aiohttp -from prometheus_client.exposition import ( - REGISTRY, - MetricsHandler, - _ThreadingSimpleServer, - choose_encoder, -) -from prometheus_client.parser import text_string_to_metric_families - -from .exceptions import InternalError - -logger = logging.getLogger(__name__) - - -class CallbackMetricsHandler(MetricsHandler): - """ - Derivative of `prometheus_client.exposition.MetricsHandler`. - - Allows callback functions to be executed when metrics are requested. - """ - - callbacks: ClassVar[list] = [] - - def do_GET(self): - """Respond to request for metrics.""" - for cb in self.callbacks: - cb() - registry = self.registry - encoder, content_type = choose_encoder(self.headers.get("Accept")) - try: - output = encoder(registry) - except Exception as err: - logger.debug(f"Error generating metric output: {err}") - self.send_error(500, "error generating metric output") - raise - self.send_response(200) - self.send_header("Content-Type", content_type) - self.end_headers() - self.wfile.write(output) - - -def start_metrics_server(port, callbacks=None, addr=""): - """Start the metrics server - - Based on `prometheus_client.exposition.start_http_server` - using a custom handler. - """ - handler = CallbackMetricsHandler.factory(REGISTRY) - if callbacks is not None: - handler.callbacks += callbacks - httpd = _ThreadingSimpleServer((addr, port), handler) - t = threading.Thread(target=httpd.serve_forever) - t.daemon = True - t.start() - - -async def get(name, port, host="127.0.0.1"): - """ - Get the value of a metric by requesting it from the prometheus web server. - - Parameters - ---------- - name : str - Name of the metric - port : int - Port of the prometheus server - host : str - Host running the prometheus server (default "127.0.0.1") - - Returns - ------- - Value of the metric - """ - async with aiohttp.ClientSession() as session: - async with session.get(f"http://{host}:{port}") as resp: - resp.raise_for_status() - metrics = await resp.text() - for family in text_string_to_metric_families(metrics): - if family.name == name: - return family.samples[0].value - raise InternalError( - f"Couldn't find metric {name} in response from coco's prometheus client." - ) diff --git a/coco/metrics.py b/coco/metrics.py new file mode 100644 index 00000000..461a25bc --- /dev/null +++ b/coco/metrics.py @@ -0,0 +1,202 @@ +"""Metrics aggregator.""" + +import json +import logging +import signal +import time +import traceback + +import redis +from prometheus_client import ( + CONTENT_TYPE_LATEST, + Counter, + Gauge, + Histogram, + disable_created_metrics, + generate_latest, +) + +logger = logging.getLogger(__name__) + +# Redis connection +redis_conn = None + +# Cache of the rendered metrics +rendered = None + + +def signal_handler(signum, frame): + """Signal handler.""" + global redis_conn + + logger.debug(f"Caught signal {signum}. Exiting.") + if redis_conn: + redis_conn.close() + redis_conn = None + raise KeyboardInterrupt + + +def init_metrics(all_endpoints): + """Initialise prometheus metric instances.""" + disable_created_metrics() + dropped_requests = Counter( + "coco_dropped_request", + "Count of requests dropped by coco.", + ["endpoint"], + ) + + # Init all the dropped requests metrics + for endpoint in all_endpoints: + dropped_requests.labels(endpoint=endpoint).reset() + + return { + "dropped_request": dropped_requests, + "coco_calls": Counter( + "coco_calls", + "Calls forwarded by coco to hosts.", + ["endpoint", "host", "port", "status"], + ), + "qlen": Gauge( + "coco_queue_length", "Length of queue storing coco requests.", unit="total" + ), + "queue_wait_time": Histogram( + "coco_queue_wait_time", + "Length of time the request is in the queue before being processed", + ["endpoint"], + unit="seconds", + ), + "external_response_time": Histogram( + "coco_external_response_time", + "Length of time external hosts take to answer coco's requests", + ["endpoint", "host", "port"], + unit="seconds", + ), + } + + +def handle_requests(pipe): + """Handle requests for metrics. + + Parameters + ---------- + pipe + The pipe to communicate with the rest of Sanic. + dirty : bool + If True, the metrics are dirty (need regeneration) + """ + global rendered + + # Are there metric requests? + if pipe.poll(): + # Re-render the metrics, if necessary + if not rendered: + rendered = (generate_latest(), CONTENT_TYPE_LATEST) + + # Once we're servicing requests, we handle them all + while pipe.poll(): + # Discard whatever we were sent + pipe.recv() + + # Reply with the render + pipe.send(rendered) + + +# qlen cache. +_qlen = None + + +def update_metrics(metrics): + """Update the coco metrics. + + If any metric is updated, the global "rendered" is reset + to force re-rendering. + + Parameters + ---------- + metrics : dict + The metrics dict created by init_metrics() + """ + global _qlen, rendered + + # Iterate over all metrics + for name, metric in metrics.items(): + if name == "dropped_request": + # Redis list "dropped_requests" has one endpoint per dropped request + while redis_conn.llen("dropped_requests") > 0: + rendered = None + metric.labels( + endpoint=redis_conn.lpop("dropped_requests").decode() + ).inc() + elif name == "coco_calls": + while redis_conn.llen("coco_calls") > 0: + rendered = None + labels = json.loads(redis_conn.lpop("coco_calls")) + metric.labels(**labels).inc() + elif name == "qlen": + new_qlen = int(redis_conn.llen("queue")) + if new_qlen != _qlen: + rendered = None + metric.set(new_qlen) + _qlen = new_qlen + elif name == "queue_wait_time" or name == "external_response_time": + while redis_conn.llen(name) > 0: + rendered = None + labels = json.loads(redis_conn.lpop(name)) + # Extract observed value from the dict + value = labels.pop("value") + metric.labels(**labels).observe(value) + else: + # Unhandled metric + logger.warning(f"not updating unknown metric {name!r}") + + +def aggregator(app, all_endpoints, pipe, port): + """Sanic worker for aggregating and reporting metrics. + + Parameters + ---------- + app: + The controlling Sanic app + all_endpoints: + A set of all endpoint names. + pipe: + The pipe on which to listen for requests. + port: + The redis port + """ + + global redis_conn + + shutdown = False + logger.debug("metric.aggregator started") + + signal.signal(signal.SIGINT, signal_handler) + signal.signal(signal.SIGTERM, signal_handler) + + try: + # Init + metrics = init_metrics(all_endpoints) + + # Connect to redis + redis_conn = redis.Redis(port=port) + + # Do two things forever + while True: + update_metrics(metrics) + handle_requests(pipe) + time.sleep(1) + except KeyboardInterrupt: + # Normal termination + logger.info("metric.aggregator shutdown") + shutdown = True + except SystemExit: + # Normal exit + logger.info("qworker shutdown") + shutdown = False + except BaseException as e: # noqa: BLE001 + logger.error(f"metric.aggregator encountered an error: {e}") + logger.error(traceback.format_exc()) + finally: + if not shutdown: + # Shutdown Sanic on abnormal exit + app.manager.terminate() diff --git a/coco/request_forwarder.py b/coco/request_forwarder.py index 6140bb10..226c98ce 100644 --- a/coco/request_forwarder.py +++ b/coco/request_forwarder.py @@ -10,10 +10,8 @@ import aiohttp import redis -from prometheus_client import Counter, Gauge, Histogram from .blocklist import Blocklist -from .metric import start_metrics_server from .result import Result from .task_pool import TaskPool from .util import Host @@ -194,19 +192,18 @@ class RequestForwarder: """ def __init__( - self, blocklist_path: os.PathLike, timeout: int, debug_connections: bool = False + self, + blocklist_path: os.PathLike, + redis_port: int, + timeout: int, + debug_connections: bool = False, ): self._endpoints = {} self._groups = {} self.session_limit = 1 self.blocklist = Blocklist([], blocklist_path) self.timeout = timeout - self.redis_conn = None - self.dropped_counter = None - self.call_counter = None - self.queue_len = None - self.queue_wait_time = None - self.response_time = None + self.redis_conn = redis.Redis(port=redis_port) self._debug_connections = debug_connections def set_session_limit(self, session_limit): @@ -250,64 +247,6 @@ def add_endpoint(self, name, endpoint): """ self._endpoints[name] = endpoint - def start_prometheus_server(self, port, redis_port): - """ - Start prometheus server. - - Parameters - ---------- - port : int - Server port. - redis_port : int - The port redis is listening on - """ - # Connect to redis - self.redis_conn = redis.Redis(host="127.0.0.1", port=redis_port, db=0) - - def fetch_request_count(): - for edpt in self._endpoints: - # Get current count and reset to 0 - incr = int(self.redis_conn.getset(f"dropped_counter_{edpt}", "0")) - self.dropped_counter.labels(endpoint=edpt).inc(incr) - - def fetch_queue_len(): - self.queue_len.set(int(self.redis_conn.llen("queue"))) - - start_metrics_server(port, callbacks=[fetch_request_count, fetch_queue_len]) - - def init_metrics(self): - """Initialise counters for every prometheus endpoint.""" - self.dropped_counter = Counter( - "coco_dropped_request", - "Count of requests dropped by coco.", - ["endpoint"], - unit="total", - ) - self.call_counter = Counter( - "coco_calls", - "Calls forwarded by coco to hosts.", - ["endpoint", "host", "port", "status"], - unit="total", - ) - self.queue_len = Gauge( - "coco_queue_length", "Length of queue storing coco requests.", unit="total" - ) - self.queue_wait_time = Histogram( - "coco_queue_wait_time", - "Length of time the request is in the queue before being processed", - ["endpoint"], - unit="seconds", - ) - self.response_time = Histogram( - "coco_external_response_time", - "Length of time external hosts take to answer coco's requests", - ["endpoint", "host", "port"], - unit="seconds", - ) - for edpt in self._endpoints: - self.dropped_counter.labels(endpoint=edpt).inc(0) - self.redis_conn.set(f"dropped_counter_{edpt}", "0") - async def internal(self, name, request=None, hosts=None, **_): """ Call an endpoint. @@ -386,12 +325,28 @@ async def _request(self, session, method, host, endpoint, request, params, timeo return host, (str(e), 0, 0) finally: response_time = time.perf_counter() - start_time - self.response_time.labels( - endpoint=endpoint, host=hostname, port=port - ).observe(response_time) - self.call_counter.labels( - endpoint=endpoint, host=hostname, port=port, status=status - ).inc() + self.redis_conn.rpush( + "external_response_time", + json.dumps( + { + "endpoint": endpoint, + "host": hostname, + "port": port, + "value": response_time, + } + ), + ) + self.redis_conn.rpush( + "coco_calls", + json.dumps( + { + "endpoint": endpoint, + "host": hostname, + "port": port, + "status": status, + } + ), + ) async def external(self, name, request, hosts, method, params=None, timeout=None): """ diff --git a/coco/worker.py b/coco/worker.py index cef77d91..4b6f66a0 100644 --- a/coco/worker.py +++ b/coco/worker.py @@ -50,11 +50,8 @@ async def _open_redis_connection(redis_port): sys.exit(1) -async def go(endpoints, redis_port, metrics_port, forwarder): +async def go(endpoints, redis_port): """Asynchronous qworker main loop.""" - # start the prometheus server for forwarded requests - forwarder.start_prometheus_server(metrics_port, redis_port) - forwarder.init_metrics() global conn, teardown conn = await _open_redis_connection(redis_port) @@ -96,7 +93,7 @@ async def go(endpoints, redis_port, metrics_port, forwarder): if received: received = float(received) queue_wait = time.perf_counter() - received - conn.rpush( + await conn.rpush( "queue_wait_time", json.dumps({"endpoint": endpoint_name, "value": queue_wait}), ) @@ -192,9 +189,7 @@ async def go(endpoints, redis_port, metrics_port, forwarder): await conn.close() -def main_loop( - app, endpoints, forwarder, coco_port, metrics_port, redis_port, frontend_timeout -): +def main_loop(app, endpoints, coco_port, redis_port, frontend_timeout): """ Wait for tasks and run them. @@ -207,12 +202,8 @@ def main_loop( endpoints : dict A dict with keys being endpoint names and values being of type :class:`Endpoint`. - forwarder - The RequestForwarder coco_port: The coco Sanic port - metrics_port: - The prometheus metrics port redis_port: The redis server port frontend_timeout : int @@ -234,9 +225,7 @@ def main_loop( scheduler = Scheduler(endpoints, "127.0.0.1", coco_port, frontend_timeout) loop.run_until_complete( - asyncio.gather( - go(endpoints, redis_port, metrics_port, forwarder), scheduler.start() - ) + asyncio.gather(go(endpoints, redis_port), scheduler.start()) ) # Cleanup diff --git a/old_tests/test_metrics.py b/old_tests/test_metrics.py deleted file mode 100644 index 7a632d3d..00000000 --- a/old_tests/test_metrics.py +++ /dev/null @@ -1,94 +0,0 @@ -"""Test prometheus metrics collection.""" - -import pytest -import requests -from coco.test import coco_runner, endpoint_farm -from prometheus_client.parser import text_string_to_metric_families - -PORT = 12056 -CONFIG = {"log_level": "INFO", "metrics_port": PORT} -ENDPT_NAME = "status" -ENDPT_NAME_FWD = "status-host" -ENDPOINTS = { - ENDPT_NAME: { - "call": {"forward": ENDPT_NAME_FWD}, - "group": "test", - "values": {"foo": "int", "bar": "str"}, - } -} -N_CALLS = 2 - - -def callback(data): - """Reply with the incoming json request.""" - return data - - -N_HOSTS = 2 -CALLBACKS = {ENDPT_NAME: callback} - - -@pytest.fixture -def farm(): - """Create a coco runner.""" - return endpoint_farm.Farm(N_HOSTS, CALLBACKS) - - -@pytest.fixture -def runner(farm): - """Create an endpoint test farm.""" - CONFIG["groups"] = {"test": farm.hosts} - with coco_runner.Runner(CONFIG, ENDPOINTS) as runner: - yield runner - - -def test_metrics(farm, runner): - """Test metrics are counting endpoint calls.""" - for i in range(N_CALLS): - runner.client(ENDPT_NAME, ["0", "1337"]) - - # Get metrics - metrics = requests.get(f"http://localhost:{PORT}/metrics") - assert metrics.status_code == 200 - metrics = text_string_to_metric_families(metrics.text) - - # parse metrics - count_coco = [] - count_forward = [] - count_wait_time = {} - for metric in metrics: - for sample in metric.samples: - if sample.name == "coco_dropped_request_total": - count_coco.append(sample) - elif sample.name == "coco_calls_total": - count_forward.append(sample) - elif sample.name == "coco_queue_wait_time_seconds_count": - count_wait_time[sample.labels["endpoint"]] = sample - - # Only expect one endpoint call - assert ( - len(count_coco) == 9 - ) # 1, plus two from internal metrics. Needs to be kept up to date - count_coco = count_coco[0] - assert list(count_coco.labels.keys()) == ["endpoint"] - assert count_coco.labels["endpoint"] == ENDPT_NAME - # No requests should have been dropped - assert count_coco.value == 0 - - # Expect one sample per host per endpoint - assert len(count_forward) == N_HOSTS - for s in count_forward: - assert set(s.labels.keys()) == {"endpoint", "host", "port", "status"} - for p in farm.ports: - ind = [int(s.labels["port"]) for s in count_forward].index(p) - assert count_forward[ind].labels["endpoint"] == ENDPT_NAME_FWD - assert count_forward[ind].labels["status"] == "200" - assert count_forward[ind].labels["host"] == "localhost" - assert count_forward[ind].value == N_CALLS - assert farm.counters()[p][ENDPT_NAME_FWD] == N_CALLS - - # Expect only queue wait time for "status", because that's all the test calls - assert len(count_wait_time) == 1 - assert count_wait_time["status"].value == 2 - - del metrics diff --git a/pyproject.toml b/pyproject.toml index 8799ada6..9aa8846e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -29,7 +29,7 @@ cocod = [ "deepdiff", "jinja2", "msgpack", - "prometheus_client<0.8", + "prometheus_client>=0.20.0", "redis>=4.2.0", "sanic>=20.6.0" ] diff --git a/tests/test_metrics.py b/tests/test_metrics.py new file mode 100644 index 00000000..a4b26a55 --- /dev/null +++ b/tests/test_metrics.py @@ -0,0 +1,64 @@ +"""Test prometheus metrics collection.""" + +from prometheus_client.parser import text_string_to_metric_families + +from coco.core import all_local_endpoints + + +def test_metrics(coco_runner): + coco_runner.add_targets("test", 2) + coco_runner.add_endpoint( + "status", + { + "call": {"forward": "status-host"}, + "group": "test", + "values": {"foo": "int", "bar": "str"}, + }, + ) + + N_CALLS = 2 + for i in range(N_CALLS): + coco_runner.client("status", "--foo=0", "--bar=1337") + + # Get metrics + result, metrics = coco_runner.call_endpoint("metrics") + assert result.status == 200 + metrics = text_string_to_metric_families(metrics) + + # parse metrics + all_dropped = [] + status_dropped = None + all_forward = [] + all_wait_time = {} + all_response_time = {} + for metric in metrics: + for sample in metric.samples: + if sample.name == "coco_dropped_request_total": + all_dropped.append(sample) + if sample.labels["endpoint"] == "status": + status_dropped = sample + elif sample.name == "coco_calls_total": + all_forward.append(sample) + elif sample.name == "coco_queue_wait_time_seconds_count": + all_wait_time[sample.labels["endpoint"]] = sample + elif sample.name == "coco_external_response_time_seconds_count": + all_response_time[sample.labels["endpoint"]] = sample + + # Expect one endpoint plus all the internal local endpoints + assert len(all_dropped) == len(all_local_endpoints) + 1 + + assert list(status_dropped.labels.keys()) == ["endpoint"] + # No requests should have been dropped + assert status_dropped.value == 0 + + # Expect one sample per host per endpoint + assert len(all_forward) == 2 + for s in all_forward: + assert set(s.labels.keys()) == {"endpoint", "host", "port", "status"} + + # Expect only queue wait time for "status", because that's all the test calls + assert len(all_wait_time) == 1 + + # Also have observations of the histograms + assert all_wait_time["status"].value == 2 + assert all_response_time["status-host"].value == 2