Skip to content

Commit 5e93f60

Browse files
authored
AAP-41776 Enable new fancy asyncio metrics for dispatcherd (#16233)
* Enable new fancy asyncio metrics for dispatcherd Remove old dispatcher metrics and patch in new data from local whatever Update test fixture to new dispatcherd version * Update dispatcherd again * Handle node filter in URL, and catch more errors * Add test for metric filter * Split module for dispatcherd metrics
1 parent 6a03115 commit 5e93f60

5 files changed

Lines changed: 112 additions & 17 deletions

File tree

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
1+
import http.client
2+
import socket
3+
import urllib.error
4+
import urllib.request
5+
import logging
6+
7+
from django.conf import settings
8+
9+
logger = logging.getLogger(__name__)
10+
11+
12+
def get_dispatcherd_metrics(request):
13+
metrics_cfg = settings.METRICS_SUBSYSTEM_CONFIG.get('server', {}).get(settings.METRICS_SERVICE_DISPATCHER, {})
14+
host = metrics_cfg.get('host', 'localhost')
15+
port = metrics_cfg.get('port', 8015)
16+
metrics_filter = []
17+
if request is not None and hasattr(request, "query_params"):
18+
try:
19+
nodes_filter = request.query_params.getlist("node")
20+
except Exception:
21+
nodes_filter = []
22+
if nodes_filter and settings.CLUSTER_HOST_ID not in nodes_filter:
23+
return ''
24+
try:
25+
metrics_filter = request.query_params.getlist("metric")
26+
except Exception:
27+
metrics_filter = []
28+
if metrics_filter:
29+
# Right now we have no way of filtering the dispatcherd metrics
30+
# so just avoid getting in the way if another metric is filtered for
31+
return ''
32+
url = f"http://{host}:{port}/metrics"
33+
try:
34+
with urllib.request.urlopen(url, timeout=1.0) as response:
35+
payload = response.read()
36+
if not payload:
37+
return ''
38+
return payload.decode('utf-8')
39+
except (urllib.error.URLError, UnicodeError, socket.timeout, TimeoutError, http.client.HTTPException) as exc:
40+
logger.debug(f"Failed to collect dispatcherd metrics from {url}: {exc}")
41+
return ''

awx/main/analytics/subsystem_metrics.py

Lines changed: 7 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
from awx.main.consumers import emit_channel_notification
1616
from awx.main.utils import is_testing
1717
from awx.main.utils.redis import get_redis_client
18+
from .dispatcherd_metrics import get_dispatcherd_metrics
1819

1920
root_key = settings.SUBSYSTEM_METRICS_REDIS_KEY_PREFIX
2021
logger = logging.getLogger('awx.main.analytics')
@@ -398,11 +399,6 @@ class DispatcherMetrics(Metrics):
398399
SetFloatM('workflow_manager_recorded_timestamp', 'Unix timestamp when metrics were last recorded'),
399400
SetFloatM('workflow_manager_spawn_workflow_graph_jobs_seconds', 'Time spent spawning workflow tasks'),
400401
SetFloatM('workflow_manager_get_tasks_seconds', 'Time spent loading workflow tasks from db'),
401-
# dispatcher subsystem metrics
402-
SetIntM('dispatcher_pool_scale_up_events', 'Number of times local dispatcher scaled up a worker since startup'),
403-
SetIntM('dispatcher_pool_active_task_count', 'Number of active tasks in the worker pool when last task was submitted'),
404-
SetIntM('dispatcher_pool_max_worker_count', 'Highest number of workers in worker pool in last collection interval, about 20s'),
405-
SetFloatM('dispatcher_availability', 'Fraction of time (in last collection interval) dispatcher was able to receive messages'),
406402
]
407403

408404
def __init__(self, *args, **kwargs):
@@ -430,8 +426,12 @@ def __init__(self, *args, **kwargs):
430426

431427
def metrics(request):
432428
output_text = ''
433-
for m in [DispatcherMetrics(), CallbackReceiverMetrics()]:
434-
output_text += m.generate_metrics(request)
429+
output_text += DispatcherMetrics().generate_metrics(request)
430+
output_text += CallbackReceiverMetrics().generate_metrics(request)
431+
432+
dispatcherd_metrics = get_dispatcherd_metrics(request)
433+
if dispatcherd_metrics:
434+
output_text += dispatcherd_metrics
435435
return output_text
436436

437437

@@ -481,13 +481,6 @@ def __init__(self):
481481
super().__init__(settings.METRICS_SERVICE_CALLBACK_RECEIVER, registry)
482482

483483

484-
class DispatcherMetricsServer(MetricsServer):
485-
def __init__(self):
486-
registry = CollectorRegistry(auto_describe=True)
487-
registry.register(CustomToPrometheusMetricsCollector(DispatcherMetrics(metrics_have_changed=False)))
488-
super().__init__(settings.METRICS_SERVICE_DISPATCHER, registry)
489-
490-
491484
class WebsocketsMetricsServer(MetricsServer):
492485
def __init__(self):
493486
registry = CollectorRegistry(auto_describe=True)

awx/main/dispatch/config.py

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,8 +38,8 @@ def get_dispatcherd_config(for_service: bool = False, mock_publish: bool = False
3838
}
3939

4040
if mock_publish:
41-
config["brokers"]["noop"] = {}
42-
config["publish"]["default_broker"] = "noop"
41+
config["brokers"]["dispatcherd.testing.brokers.noop"] = {}
42+
config["publish"]["default_broker"] = "dispatcherd.testing.brokers.noop"
4343
else:
4444
config["brokers"]["pg_notify"] = {
4545
"config": get_pg_notify_params(),
@@ -56,5 +56,11 @@ def get_dispatcherd_config(for_service: bool = False, mock_publish: bool = False
5656
}
5757

5858
config["brokers"]["pg_notify"]["channels"] = ['tower_broadcast_all', 'tower_settings_change', get_task_queuename()]
59+
metrics_cfg = settings.METRICS_SUBSYSTEM_CONFIG.get('server', {}).get(settings.METRICS_SERVICE_DISPATCHER)
60+
if metrics_cfg:
61+
config["service"]["metrics_kwargs"] = {
62+
"host": metrics_cfg.get("host", "localhost"),
63+
"port": metrics_cfg.get("port", 8015),
64+
}
5965

6066
return config

awx/main/tests/functional/analytics/test_metrics.py

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,11 @@
11
import pytest
22

3+
from django.test import RequestFactory
34
from prometheus_client.parser import text_string_to_metric_families
5+
from rest_framework.request import Request
46
from awx.main import models
57
from awx.main.analytics.metrics import metrics
8+
from awx.main.analytics.dispatcherd_metrics import get_dispatcherd_metrics
69
from awx.api.versioning import reverse
710

811
EXPECTED_VALUES = {
@@ -77,3 +80,55 @@ def test_metrics_http_methods(get, post, patch, put, options, admin):
7780
assert patch(get_metrics_view_db_only(), user=admin).status_code == 405
7881
assert post(get_metrics_view_db_only(), user=admin).status_code == 405
7982
assert options(get_metrics_view_db_only(), user=admin).status_code == 200
83+
84+
85+
class DummyMetricsResponse:
86+
def __init__(self, payload):
87+
self._payload = payload
88+
89+
def read(self):
90+
return self._payload
91+
92+
def __enter__(self):
93+
return self
94+
95+
def __exit__(self, exc_type, exc, tb):
96+
return False
97+
98+
99+
def test_dispatcherd_metrics_node_filter_match(mocker, settings):
100+
settings.CLUSTER_HOST_ID = "awx-1"
101+
payload = b'# HELP test_metric A test metric\n# TYPE test_metric gauge\ntest_metric 1\n'
102+
103+
def fake_urlopen(url, timeout=1.0):
104+
return DummyMetricsResponse(payload)
105+
106+
mocker.patch('urllib.request.urlopen', fake_urlopen)
107+
108+
request = Request(RequestFactory().get('/api/v2/metrics/', {'node': 'awx-1'}))
109+
110+
assert get_dispatcherd_metrics(request) == payload.decode('utf-8')
111+
112+
113+
def test_dispatcherd_metrics_node_filter_excludes_local(mocker, settings):
114+
settings.CLUSTER_HOST_ID = "awx-1"
115+
116+
def fake_urlopen(*args, **kwargs):
117+
raise AssertionError("urlopen should not be called when node filter excludes local node")
118+
119+
mocker.patch('urllib.request.urlopen', fake_urlopen)
120+
121+
request = Request(RequestFactory().get('/api/v2/metrics/', {'node': 'awx-2'}))
122+
123+
assert get_dispatcherd_metrics(request) == ''
124+
125+
126+
def test_dispatcherd_metrics_metric_filter_excludes_unrelated(mocker):
127+
def fake_urlopen(*args, **kwargs):
128+
raise AssertionError("urlopen should not be called when metric filter excludes dispatcherd metrics")
129+
130+
mocker.patch('urllib.request.urlopen', fake_urlopen)
131+
132+
request = Request(RequestFactory().get('/api/v2/metrics/', {'metric': 'awx_system_info'}))
133+
134+
assert get_dispatcherd_metrics(request) == ''

requirements/requirements.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,7 @@ cython==3.1.3
116116
# via -r /awx_devel/requirements/requirements.in
117117
daphne==4.2.1
118118
# via -r /awx_devel/requirements/requirements.in
119-
dispatcherd[pg-notify]==2025.12.12
119+
dispatcherd[pg-notify]==2026.01.27
120120
# via -r /awx_devel/requirements/requirements.in
121121
distro==1.9.0
122122
# via -r /awx_devel/requirements/requirements.in

0 commit comments

Comments
 (0)