@cryptotaxi247 / infra / commits / 786fb960

prometheus: drop queue-runner reexporter

Not relevant any more wit the new queue-runner.

Martin Weinelt committed Jun 18, 2026 at 00:40 UTC 786fb9607c463eae7886e804d01d19a70f12213e
2 files changed -613
build/pluto/prometheus/exporters/hydra-queue-runner-reexporter.py deleted
-583
@@ -1,583 +0,0 @@
1 -#!/usr/bin/env nix-shell
2 -#!nix-shell -i python3 -p python3 -p python3Packages.requests -p python3Packages.prometheus_client
3 -
4 -import contextlib
5 -import json
6 -import time
7 -
8 -import requests
9 -from prometheus_client import CollectorRegistry, start_http_server
10 -from prometheus_client.core import CounterMetricFamily, GaugeMetricFamily
11 -
12 -
13 -def debug_remaining_state(edict) -> None:
14 - # pprint(edict.remaining_state())
15 - pass
16 -
17 -
18 -class EvaporatingDict:
19 - def __init__(self, state) -> None:
20 - self._state = state
21 -
22 - def preserving_read(self, key):
23 - val = self._state[key]
24 -
25 - if isinstance(val, dict):
26 - return EvaporatingDict(val)
27 - return val
28 -
29 - def preserving_read_default(self, key, default):
30 - try:
31 - return self.preserving_read(key)
32 - except KeyError:
33 - return default
34 -
35 - def destructive_read(self, key):
36 - val = self.preserving_read(key)
37 - del self._state[key]
38 - return val
39 -
40 - def destructive_read_default(self, key, default):
41 - try:
42 - val = self.preserving_read(key)
43 - del self._state[key]
44 - return val
45 - except KeyError:
46 - # Not nice, but accounts for weird conditionals in Hydra
47 - # todo: log bad reads?
48 - return default
49 -
50 - def unused_read(self, key) -> None:
51 - self.destructive_read_default(key, default=None)
52 -
53 - def remaining_state(self):
54 - return self._state
55 -
56 - def items(self):
57 - keys = list(self._state.keys())
58 - for key in keys:
59 - yield (key, self.destructive_read(key))
60 -
61 -
62 -class HydraScrapeImporter:
63 - def __init__(self, status) -> None:
64 - self._status = EvaporatingDict(status)
65 -
66 - def collect(self):
67 - # The metrics are consumed in the order presented by
68 - # https://github.com/NixOS/hydra/blob/adf59a395993d5ed1d7a31108f7666195f789c99/src/hydra-queue-runner/hydra-queue-runner.cc#L536
69 - yield self.trivial_gauge(
70 - "up",
71 - "Is hydra running",
72 - 1 if self.destructive_read("status") == "up" else 0,
73 - )
74 - yield self.trivial_counter(
75 - "time", "Hydra's current time", self.destructive_read("time")
76 - )
77 - yield self.trivial_counter(
78 - "uptime", "Hydra's uptime", self.destructive_read("uptime")
79 - )
80 - self.unused_metric("pid")
81 - yield self.trivial_gauge(
82 - "builds_queued",
83 - "Current build queue size",
84 - self.destructive_read("nrQueuedBuilds"),
85 - )
86 - yield self.trivial_gauge(
87 - "steps_queued",
88 - "Current number of steps for the build queue",
89 - self.destructive_read("nrUnfinishedSteps"),
90 - )
91 - yield self.trivial_gauge(
92 - "steps_runnable",
93 - "Current number of steps which can run immediately",
94 - self.destructive_read("nrRunnableSteps"),
95 - )
96 - yield self.trivial_gauge(
97 - "steps_active",
98 - "Current number of steps which are currently active",
99 - self.destructive_read("nrActiveSteps"),
100 - )
101 - yield self.trivial_gauge(
102 - "steps_building",
103 - "Current number of steps which are currently building",
104 - self.destructive_read("nrStepsBuilding"),
105 - )
106 - yield self.trivial_gauge(
107 - "steps_copying_to",
108 - "Current number of steps which are having build inputs copied to a builder",
109 - self.destructive_read("nrStepsCopyingTo"),
110 - )
111 - yield self.trivial_gauge(
112 - "steps_copying_from",
113 - "Current number of steps which are having build results copied from a builder",
114 - self.destructive_read("nrStepsCopyingFrom"),
115 - )
116 - yield self.trivial_gauge(
117 - "steps_waiting",
118 - "Current number of steps which are waiting",
119 - self.destructive_read("nrStepsWaiting"),
120 - )
121 - yield self.trivial_counter(
122 - "build_inputs_sent_bytes",
123 - "Total count of bytes sent due to build inputs",
124 - self.destructive_read("bytesSent"),
125 - )
126 - yield self.trivial_counter(
127 - "build_outputs_received_bytes",
128 - "Total count of bytes received from build outputs",
129 - self.destructive_read("bytesReceived"),
130 - )
131 - yield self.trivial_counter(
132 - "builds_read",
133 - "Total count of builds whose outputs have been read",
134 - self.destructive_read("nrBuildsRead"),
135 - )
136 - yield self.trivial_counter(
137 - "builds_read_seconds",
138 - "Total number of seconds spent reading build outputs",
139 - self.destructive_read("buildReadTimeMs") / 1000,
140 - )
141 - self.unused_metric("buildReadTimeAvgMs") # implementable in prometheus queries
142 -
143 - yield self.trivial_counter(
144 - "builds_done",
145 - "Total count of builds performed",
146 - self.destructive_read("nrBuildsDone"),
147 - )
148 - yield self.trivial_counter(
149 - "steps_started",
150 - "Total count of steps started",
151 - self.destructive_read("nrStepsStarted"),
152 - )
153 - yield self.trivial_counter(
154 - "steps_done",
155 - "Total count of steps completed",
156 - self.destructive_read("nrStepsDone"),
157 - )
158 - yield self.trivial_counter(
159 - "retries", "Total count of retries", self.destructive_read("nrRetries")
160 - )
161 - yield self.trivial_counter(
162 - "max_retries",
163 - "Maximum count of retries for any single job",
164 - self.destructive_read("maxNrRetries"),
165 - )
166 - yield self.trivial_counter(
167 - "step_time",
168 - "Total time spent executing steps",
169 - self.destructive_read_default("totalStepTime", 0),
170 - )
171 - yield self.trivial_counter(
172 - "step_build_time",
173 - "Total time spent executing builds steps (???)",
174 - self.destructive_read_default("totalStepBuildTime", 0),
175 - )
176 - self.unused_metric("avgStepTime")
177 - self.unused_metric("avgStepBuildTime")
178 -
179 - yield self.trivial_counter(
180 - "queue_wakeup",
181 - "Count of the times the queue runner has been notified of queue changes",
182 - self.destructive_read("nrQueueWakeups"),
183 - )
184 - yield self.trivial_counter(
185 - "dispatcher_wakeup",
186 - "Count of the times the queue runner work dispatcher woke up due to new runnable builds and completed builds.",
187 - self.destructive_read("nrDispatcherWakeups"),
188 - )
189 - yield self.trivial_counter(
190 - "dispatch_execution_seconds",
191 - "Number of seconds the dispatcher has spent working",
192 - self.destructive_read("dispatchTimeMs") / 1000,
193 - )
194 - self.unused_metric("dispatchTimeAvgMs")
195 -
196 - yield self.trivial_gauge(
197 - "db_connections",
198 - "Number of connections to the database",
199 - self.destructive_read("nrDbConnections"),
200 - )
201 - yield self.trivial_gauge(
202 - "db_updates",
203 - "Number of in-progress database updates",
204 - self.destructive_read("nrActiveDbUpdates"),
205 - )
206 - yield self.trivial_counter(
207 - "notifications_total",
208 - "Total number of notifications sent",
209 - self.preserving_read_default("nrNotificationsDone", 0)
210 - + self.preserving_read_default("nrNotificationsFailed", 0),
211 - )
212 - yield self.trivial_counter(
213 - "notifications_done",
214 - "Number of notifications completed",
215 - self.destructive_read_default("nrNotificationsDone", 0),
216 - )
217 - yield self.trivial_counter(
218 - "notifications_failed",
219 - "Number of notifications failed",
220 - self.destructive_read_default("nrNotificationsFailed", 0),
221 - )
222 - yield self.trivial_counter(
223 - "notifications_in_progress",
224 - "Number of notifications in_progress",
225 - self.destructive_read_default("nrNotificationsInProgress", 0),
226 - )
227 - yield self.trivial_counter(
228 - "notifications_pending",
229 - "Number of notifications pending",
230 - self.destructive_read_default("nrNotificationsPending", 0),
231 - )
232 - yield self.trivial_counter(
233 - "notifications_seconds",
234 - "Time spent delivering notifications",
235 - self.destructive_read_default("nrNotificationTimeMs", 0) / 1000,
236 - )
237 - self.unused_metric("nrNotificationTimeAvgMs")
238 -
239 - machineCollector = MachineScrapeImporter()
240 - for name, report in self.destructive_read("machines").items():
241 - machineCollector.load_machine(name, report)
242 - for metric in machineCollector.metrics():
243 - yield metric
244 -
245 - jobsetCollector = JobsetScrapeImporter()
246 - for name, report in self.destructive_read("jobsets").items():
247 - jobsetCollector.load_jobset(name, report)
248 - for metric in jobsetCollector.metrics():
249 - yield metric
250 -
251 - machineTypesCollector = MachineTypeScrapeImporter()
252 - for name, report in self.destructive_read("machineTypes").items():
253 - machineTypesCollector.load_machine_type(name, report)
254 - for metric in machineTypesCollector.metrics():
255 - yield metric
256 -
257 - store = self.destructive_read("store")
258 - yield self.trivial_counter(
259 - "store_nar_info_read",
260 - "Number of NarInfo files read from the binary cache",
261 - store.destructive_read("narInfoRead"),
262 - )
263 - yield self.trivial_counter(
264 - "store_nar_info_read_averted",
265 - "Number of NarInfo files reads which were avoided",
266 - store.destructive_read("narInfoReadAverted"),
267 - )
268 - yield self.trivial_counter(
269 - "store_nar_info_missing",
270 - "Number of NarInfo files read attempts which identified a missing narinfo file",
271 - store.destructive_read("narInfoMissing"),
272 - )
273 - yield self.trivial_counter(
274 - "store_nar_info_write",
275 - "Number of NarInfo files written to the binary cache",
276 - store.destructive_read("narInfoWrite"),
277 - )
278 - yield self.trivial_gauge(
279 - "store_nar_info_cache_size",
280 - "Size of the in-memory store path information cache",
281 - store.destructive_read("narInfoCacheSize"),
282 - )
283 - yield self.trivial_counter(
284 - "store_nar_read",
285 - "Number of NAR files read from the binary cache",
286 - store.destructive_read("narRead"),
287 - )
288 - yield self.trivial_counter(
289 - "store_nar_read_bytes",
290 - "Number of NAR file bytes read after decompression from the binary cache",
291 - store.destructive_read("narReadBytes"),
292 - )
293 - yield self.trivial_counter(
294 - "store_nar_read_compressed_bytes",
295 - "Number of NAR file bytes read before decompression from the binary cache",
296 - store.destructive_read("narReadCompressedBytes"),
297 - )
298 - yield self.trivial_counter(
299 - "store_nar_write",
300 - "Number of NAR files written to the binary cache",
301 - store.destructive_read("narWrite"),
302 - )
303 - yield self.trivial_counter(
304 - "store_nar_write_averted",
305 - "Number of NAR files writes skipped due to the NAR already being in the binary cache",
306 - store.destructive_read("narWriteAverted"),
307 - )
308 - yield self.trivial_counter(
309 - "store_nar_write_bytes",
310 - "Number of NAR file bytes written after decompression to the binary cache",
311 - store.destructive_read("narWriteBytes"),
312 - )
313 - yield self.trivial_counter(
314 - "store_nar_write_compressed_bytes",
315 - "Number of NAR file bytes written before decompression to the binary cache",
316 - store.destructive_read("narWriteCompressedBytes"),
317 - )
318 - yield self.trivial_counter(
319 - "store_nar_write_compression_seconds",
320 - "Number of seconds spent compressing data when writing NARs to the binary cache",
321 - store.destructive_read("narWriteCompressionTimeMs") / 1000,
322 - )
323 - store.unused_read("narCompressionSavings")
324 - store.unused_read("narCompressionSpeed")
325 -
326 - try:
327 - s3 = self.destructive_read("s3")
328 - except KeyError:
329 - # no key, no metrics
330 - s3 = None
331 - if s3:
332 - # Not in the above try to avoid the try catching mistakes
333 - # in the following code
334 - yield self.trivial_counter(
335 - "store_s3_put", "Number of PUTs to S3", s3.destructive_read("put")
336 - )
337 - yield self.trivial_counter(
338 - "store_s3_put_bytes",
339 - "Number of bytes written to S3",
340 - s3.destructive_read("putBytes"),
341 - )
342 - yield self.trivial_counter(
343 - "store_s3_put_seconds",
344 - "Number of seconds spent writing to S3",
345 - s3.destructive_read("putTimeMs") / 1000,
346 - )
347 - s3.unused_read("putSpeed")
348 - yield self.trivial_counter(
349 - "store_s3_get", "Number of GETs to S3", s3.destructive_read("get")
350 - )
351 - yield self.trivial_counter(
352 - "store_s3_get_bytes",
353 - "Number of bytes read from S3",
354 - s3.destructive_read("getBytes"),
355 - )
356 - yield self.trivial_counter(
357 - "store_s3_get_seconds",
358 - "Number of seconds spent reading from S3",
359 - s3.destructive_read("getTimeMs") / 1000,
360 - )
361 - s3.unused_read("getSpeed")
362 -
363 - yield self.trivial_counter(
364 - "store_s3_head", "Number of HEADs to S3", s3.destructive_read("head")
365 - )
366 - yield self.trivial_counter(
367 - "store_s3_cost_approximate_dollars",
368 - "Estimated cost of the S3 bucket activity",
369 - s3.destructive_read("costDollarApprox"),
370 - )
371 - debug_remaining_state(s3)
372 - debug_remaining_state(store)
373 -
374 - def trivial_gauge(self, name, help, value):
375 - c = GaugeMetricFamily(f"hydra_{name}", help)
376 - c.add_metric([], value)
377 - return c
378 -
379 - def trivial_counter(self, name, help, value):
380 - c = CounterMetricFamily(f"hydra_{name}_total", help)
381 - c.add_metric([], value)
382 - return c
383 -
384 - def unused_metric(self, key) -> None:
385 - self._status.unused_read(key)
386 -
387 - def preserving_read(self, key):
388 - return self._status.preserving_read(key)
389 -
390 - def preserving_read_default(self, key, default):
391 - return self._status.preserving_read_default(key, default)
392 -
393 - def destructive_read(self, key):
394 - return self._status.destructive_read(key)
395 -
396 - def destructive_read_default(self, key, default):
397 - return self._status.destructive_read_default(key, default)
398 -
399 - def uncollected_status(self):
400 - return self._status.remaining_state()
401 -
402 -
403 -def blackhole(*args, **kwargs) -> None:
404 - return None
405 -
406 -
407 -class MachineScrapeImporter:
408 - def __init__(self) -> None:
409 - labels = ["host"]
410 - self.consective_failures = GaugeMetricFamily(
411 - "hydra_machine_consecutive_failures",
412 - "Number of consecutive failed builds",
413 - labels=labels,
414 - )
415 - self.current_jobs = GaugeMetricFamily(
416 - "hydra_machine_current_jobs", "Number of current jobs", labels=labels
417 - )
418 - self.idle_since = GaugeMetricFamily(
419 - "hydra_machine_idle_since",
420 - "When the current idle period started",
421 - labels=labels,
422 - )
423 - self.disabled_until = GaugeMetricFamily(
424 - "hydra_machine_disabled_until",
425 - "When the machine will be used again",
426 - labels=labels,
427 - )
428 - self.enabled = GaugeMetricFamily(
429 - "hydra_machine_enabled",
430 - "If the machine is enabled (1) or not (0)",
431 - labels=labels,
432 - )
433 - self.last_failure = CounterMetricFamily(
434 - "hydra_machine_last_failure", "timestamp of the last failure", labels=labels
435 - )
436 - self.number_steps_done = CounterMetricFamily(
437 - "hydra_machine_steps_done_total",
438 - "Total count of the steps completed",
439 - labels=labels,
440 - )
441 - self.total_step_build_time = CounterMetricFamily(
442 - "hydra_machine_step_build_time_total",
443 - "Number of seconds spent building steps",
444 - labels=labels,
445 - )
446 - self.total_step_time = CounterMetricFamily(
447 - "hydra_machine_step_time_total",
448 - "Number of seconds spent on steps",
449 - labels=labels,
450 - )
451 -
452 - def load_machine(self, name, report) -> None:
453 - report.unused_read("mandatoryFeatures")
454 - report.unused_read("supportedFeatures")
455 - report.unused_read("systemTypes")
456 - report.unused_read("avgStepBuildTime")
457 - report.unused_read("avgStepTime")
458 - labels = [name]
459 - self.consective_failures.add_metric(
460 - labels, report.destructive_read("consecutiveFailures")
461 - )
462 - self.current_jobs.add_metric(labels, report.destructive_read("currentJobs"))
463 - with contextlib.suppress(KeyError):
464 - self.idle_since.add_metric(labels, report.destructive_read("idleSince"))
465 - self.disabled_until.add_metric(labels, report.destructive_read("disabledUntil"))
466 - self.enabled.add_metric(labels, 1 if report.destructive_read("enabled") else 0)
467 - self.last_failure.add_metric(labels, report.destructive_read("lastFailure"))
468 - self.number_steps_done.add_metric(
469 - labels, report.destructive_read("nrStepsDone")
470 - )
471 - self.total_step_build_time.add_metric(
472 - labels, report.destructive_read_default("totalStepBuildTime", default=0)
473 - )
474 - self.total_step_time.add_metric(
475 - labels, report.destructive_read_default("totalStepTime", default=0)
476 - )
477 - debug_remaining_state(report)
478 -
479 - def metrics(self):
480 - yield self.consective_failures
481 - yield self.current_jobs
482 - yield self.idle_since
483 - yield self.disabled_until
484 - yield self.enabled
485 - yield self.last_failure
486 - yield self.number_steps_done
487 - yield self.total_step_build_time
488 - yield self.total_step_time
489 -
490 -
491 -class JobsetScrapeImporter:
492 - def __init__(self) -> None:
493 - self.seconds = CounterMetricFamily(
494 - "hydra_jobset_seconds_total",
495 - "Total number of seconds the jobset has been building",
496 - labels=["name"],
497 - )
498 - self.shares_used = CounterMetricFamily(
499 - "hydra_jobset_shares_used_total",
500 - "Total shares the jobset has consumed",
501 - labels=["name"],
502 - )
503 -
504 - def load_jobset(self, name, report) -> None:
505 - self.seconds.add_metric([name], report.destructive_read("seconds"))
506 - self.shares_used.add_metric([name], report.destructive_read("shareUsed"))
507 - debug_remaining_state(report)
508 -
509 - def metrics(self):
510 - yield self.seconds
511 - yield self.shares_used
512 -
513 -
514 -class MachineTypeScrapeImporter:
515 - def __init__(self) -> None:
516 - self.runnable = GaugeMetricFamily(
517 - "hydra_machine_type_runnable",
518 - "Number of currently runnable builds",
519 - labels=["machineType"],
520 - )
521 - self.running = GaugeMetricFamily(
522 - "hydra_machine_type_running",
523 - "Number of currently running builds",
524 - labels=["machineType"],
525 - )
526 - self.wait_time = CounterMetricFamily(
527 - "hydra_machine_type_wait_time_total",
528 - "Number of seconds spent waiting",
529 - labels=["machineType"],
530 - )
531 - self.last_active = CounterMetricFamily(
532 - "hydra_machine_type_last_active_total",
533 - "Last time this machine type was active",
534 - labels=["machineType"],
535 - )
536 -
537 - def load_machine_type(self, name, report) -> None:
538 - self.runnable.add_metric([name], report.destructive_read("runnable"))
539 - self.running.add_metric([name], report.destructive_read("running"))
540 - with contextlib.suppress(KeyError):
541 - self.wait_time.add_metric([name], report.destructive_read("waitTime"))
542 - with contextlib.suppress(KeyError):
543 - self.last_active.add_metric([name], report.destructive_read("lastActive"))
544 -
545 - debug_remaining_state(report)
546 -
547 - def metrics(self):
548 - yield self.runnable
549 - yield self.running
550 - yield self.wait_time
551 - yield self.last_active
552 -
553 -
554 -class ScrapeCollector:
555 - def __init__(self) -> None:
556 - pass
557 -
558 - def collect(self):
559 - return HydraScrapeImporter(scrape()).collect()
560 -
561 -
562 -def scrape(cached=None):
563 - if cached:
564 - with open(cached) as f:
565 - return json.load(f)
566 - else:
567 - print("Scraping")
568 - return requests.get(
569 - "https://hydra.nixos.org/queue-runner-status",
570 - headers={"Content-Type": "application/json"},
571 - ).json()
572 -
573 -
574 -registry = CollectorRegistry()
575 -
576 -registry.register(ScrapeCollector())
577 -
578 -if __name__ == "__main__":
579 - # Start up the server to expose the metrics.
580 - start_http_server(9200, registry=registry)
581 - # Generate some requests.
582 - while True:
583 - time.sleep(30)
build/pluto/prometheus/exporters/hydra.nix
-30
@@ -1,31 +1,6 @@
1 { pkgs, ... }:
2
3 {
4 - systemd.services.prometheus-hydra-queue-runner-exporter = {
5 - wantedBy = [ "multi-user.target" ];
6 - after = [ "network.target" ];
7 - wants = [ "network.target" ];
8 - serviceConfig = {
9 - DynamicUser = true;
10 - Restart = "always";
11 - RestartSec = "60s";
12 - PrivateTmp = true;
13 - WorkingDirectory = "/tmp";
14 - ExecStart =
15 - let
16 - python = pkgs.python3.withPackages (
17 - ps: with ps; [
18 - requests
19 - prometheus-client
20 - ]
21 - );
22 - in
23 - ''
24 - ${python.interpreter} ${./hydra-queue-runner-reexporter.py}
25 - '';
26 - };
27 - };
28 -
4 services.prometheus = {
5 scrapeConfigs = [
6 {
@@ -52,11 +27,6 @@
27 scheme = "https";
28 static_configs = [ { targets = [ "hydra.nixos.org:443" ]; } ];
29 }
55 - {
56 - job_name = "hydra-reexport";
57 - metrics_path = "/";
58 - static_configs = [ { targets = [ "localhost:9200" ]; } ];
59 - }
30 ];
31
32 ruleFiles = [