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