main
py 513 lines 20.8 KB
Raw
1 import math
2 from typing import Any
3
4 from loguru import logger
5 from sqlalchemy.ext.asyncio import AsyncSession
6
7 from app.connectors.influxdb.schema.metrics import HostsResponse
8 from app.connectors.influxdb.schema.metrics import MetricsResponse
9 from app.connectors.influxdb.utils.universal import create_influxdb_client
10 from app.connectors.influxdb.utils.universal import get_influxdb_organization
11 from app.connectors.utils import get_connector_info_from_db
12
13 RANGE_WINDOWS = {
14 "1": "1m",
15 "3": "2m",
16 "6": "5m",
17 "12": "10m",
18 "24": "20m",
19 "48": "30m",
20 "72": "1h",
21 "168": "2h",
22 "720": "12h",
23 }
24
25
26 async def _get_influxdb_bucket(session: AsyncSession) -> str:
27 """
28 Read the `connector_extra_data` from the database and return the bucket name.
29 The extra data format is: `ORG,BUCKET` e.g. `SOCFORTRESS,telegraf`.
30 """
31 attributes = await get_connector_info_from_db("InfluxDB", session)
32 if attributes is None:
33 raise ValueError("No InfluxDB connector found in the database")
34 parts = attributes["connector_extra_data"].split(",")
35 if len(parts) < 2:
36 raise ValueError(
37 "Invalid connector_extra_data format for InfluxDB. Expected 'ORG,BUCKET'.",
38 )
39 return parts[1].strip()
40
41
42 def _window(range_h: str) -> str:
43 """Return an appropriate aggregateWindow interval for the given hour range."""
44 return RANGE_WINDOWS.get(range_h, "5m")
45
46
47 def _range_clause(range_h: str) -> str:
48 return f"range(start: -{range_h}h)"
49
50
51 def _parse_ts_series(
52 result,
53 label_field: str = "_field",
54 ) -> dict[str, list[dict[str, Any]]]:
55 """
56 Convert InfluxDB table results into a dict of named series.
57 Each series is a list of {time, value} dicts.
58 """
59 series: dict[str, list[dict[str, Any]]] = {}
60 for table in result:
61 for record in table.records:
62 label = record.values.get(label_field, record.values.get("_field", "unknown"))
63 point = {
64 "time": record.get_time().isoformat(),
65 "value": record.get_value(),
66 }
67 series.setdefault(label, []).append(point)
68 return series
69
70
71 def _last_value(result) -> Any:
72 """Return the scalar _value from the last record, or None."""
73 for table in result:
74 for record in table.records:
75 return record.get_value()
76 return None
77
78
79 async def _run_queries(
80 influxdb_client,
81 org: str,
82 queries: dict[str, str],
83 ts_keys: set[str] | None = None,
84 label_field: str = "_field",
85 ) -> dict[str, Any]:
86 """
87 Execute a batch of Flux queries and return parsed results.
88 Keys in *ts_keys* are parsed as time-series; others as last-value scalars.
89 """
90 ts_keys = ts_keys or set()
91 query_api = influxdb_client.query_api()
92 results: dict[str, Any] = {}
93 for key, flux in queries.items():
94 result = await query_api.query(flux, org=org)
95 if key in ts_keys:
96 results[key] = _parse_ts_series(result, label_field=label_field)
97 else:
98 results[key] = _last_value(result)
99 return results
100
101
102 # ── Hosts ────────────────────────────────────────────────────────────────
103
104
105 async def get_hosts(session: AsyncSession) -> HostsResponse:
106 """Retrieve the list of unique host tag values from the metrics bucket."""
107 connector_info = await get_connector_info_from_db("InfluxDB", session)
108 if not connector_info:
109 return HostsResponse(success=False, message="InfluxDB connector not found", hosts=[])
110
111 try:
112 bucket = await _get_influxdb_bucket(session)
113 org = await get_influxdb_organization()
114 except Exception as e:
115 logger.error(f"Error getting InfluxDB config: {e}")
116 return HostsResponse(success=False, message=str(e), hosts=[])
117
118 influxdb_client = await create_influxdb_client("InfluxDB")
119 try:
120 flux = f'import "influxdata/influxdb/schema"\n' f'schema.tagValues(bucket: "{bucket}", tag: "host")'
121 query_api = influxdb_client.query_api()
122 result = await query_api.query(flux, org=org)
123 hosts = sorted(
124 {record.get_value() for table in result for record in table.records if record.get_value()},
125 )
126 return HostsResponse(success=True, message="Successfully retrieved hosts", hosts=hosts)
127 except Exception as e:
128 logger.error(f"Error fetching hosts: {e}")
129 return HostsResponse(success=False, message=f"Error fetching hosts: {e}", hosts=[])
130 finally:
131 await influxdb_client.close()
132
133
134 # ── Summary ──────────────────────────────────────────────────────────────
135
136
137 async def get_summary(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
138 """Retrieve a summary of system metrics for a given host."""
139 try:
140 bucket = await _get_influxdb_bucket(session)
141 org = await get_influxdb_organization()
142 except Exception as e:
143 return MetricsResponse(success=False, message=str(e))
144
145 influxdb_client = await create_influxdb_client("InfluxDB")
146 rng = _range_clause(range_h)
147 win = _window(range_h)
148
149 queries = {
150 "uptime": f"""from(bucket: "{bucket}") |> {rng}
151 |> filter(fn: (r) => r["_measurement"] == "system")
152 |> filter(fn: (r) => r["_field"] == "uptime")
153 |> filter(fn: (r) => r["host"] == "{host}")
154 |> last()""",
155 "total_mem": f"""from(bucket: "{bucket}") |> {rng}
156 |> filter(fn: (r) => r["_measurement"] == "mem")
157 |> filter(fn: (r) => r["_field"] == "total")
158 |> filter(fn: (r) => r["host"] == "{host}")
159 |> last()""",
160 "cpus": f"""from(bucket: "{bucket}") |> {rng}
161 |> filter(fn: (r) => r["_measurement"] == "system")
162 |> filter(fn: (r) => r["_field"] == "n_cpus")
163 |> filter(fn: (r) => r["host"] == "{host}")
164 |> last()""",
165 "total_processes": f"""from(bucket: "{bucket}") |> {rng}
166 |> filter(fn: (r) => r["_measurement"] == "processes")
167 |> filter(fn: (r) => r["_field"] == "total")
168 |> filter(fn: (r) => r["host"] == "{host}")
169 |> last()""",
170 "cpu_idle": f"""from(bucket: "{bucket}") |> {rng}
171 |> filter(fn: (r) => r["_measurement"] == "cpu")
172 |> filter(fn: (r) => r["_field"] == "usage_idle")
173 |> filter(fn: (r) => r["cpu"] == "cpu-total")
174 |> filter(fn: (r) => r["host"] == "{host}")
175 |> last()""",
176 "logged_on_users": f"""from(bucket: "{bucket}") |> {rng}
177 |> filter(fn: (r) => r["_measurement"] == "system")
178 |> filter(fn: (r) => r["_field"] == "n_users")
179 |> filter(fn: (r) => r["host"] == "{host}")
180 |> last()""",
181 "swap_free": f"""from(bucket: "{bucket}") |> {rng}
182 |> filter(fn: (r) => r["_measurement"] == "swap")
183 |> filter(fn: (r) => r["_field"] == "free")
184 |> filter(fn: (r) => r["host"] == "{host}")
185 |> last()""",
186 "load": f"""from(bucket: "{bucket}") |> {rng}
187 |> filter(fn: (r) => r["_measurement"] == "system")
188 |> filter(fn: (r) => r["_field"] == "load1" or r["_field"] == "load5" or r["_field"] == "load15")
189 |> filter(fn: (r) => r["host"] == "{host}")
190 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
191 }
192
193 try:
194 data = await _run_queries(influxdb_client, org, queries, ts_keys={"load"})
195 return MetricsResponse(success=True, message="Successfully retrieved summary", data=data)
196 except Exception as e:
197 logger.error(f"Error fetching summary: {e}")
198 return MetricsResponse(success=False, message=f"Error fetching summary: {e}")
199 finally:
200 await influxdb_client.close()
201
202
203 # ── CPU ──────────────────────────────────────────────────────────────────
204
205
206 async def get_cpu_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
207 """Retrieve CPU time-series metrics for a given host."""
208 try:
209 bucket = await _get_influxdb_bucket(session)
210 org = await get_influxdb_organization()
211 except Exception as e:
212 return MetricsResponse(success=False, message=str(e))
213
214 influxdb_client = await create_influxdb_client("InfluxDB")
215 rng = _range_clause(range_h)
216 win = _window(range_h)
217
218 queries = {
219 "cpu_usage_system": f"""from(bucket: "{bucket}") |> {rng}
220 |> filter(fn: (r) => r["_measurement"] == "cpu")
221 |> filter(fn: (r) => r["_field"] == "usage_system")
222 |> filter(fn: (r) => r["cpu"] == "cpu-total")
223 |> filter(fn: (r) => r["host"] == "{host}")
224 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
225 "cpu_usage_user": f"""from(bucket: "{bucket}") |> {rng}
226 |> filter(fn: (r) => r["_measurement"] == "cpu")
227 |> filter(fn: (r) => r["_field"] == "usage_user")
228 |> filter(fn: (r) => r["cpu"] == "cpu-total")
229 |> filter(fn: (r) => r["host"] == "{host}")
230 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
231 "cpu_iowait": f"""from(bucket: "{bucket}") |> {rng}
232 |> filter(fn: (r) => r["_measurement"] == "cpu")
233 |> filter(fn: (r) => r["_field"] == "usage_iowait")
234 |> filter(fn: (r) => r["cpu"] == "cpu-total")
235 |> filter(fn: (r) => r["host"] == "{host}")
236 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
237 "cpu_softirq": f"""from(bucket: "{bucket}") |> {rng}
238 |> filter(fn: (r) => r["_measurement"] == "cpu")
239 |> filter(fn: (r) => r["_field"] == "usage_softirq")
240 |> filter(fn: (r) => r["cpu"] == "cpu-total")
241 |> filter(fn: (r) => r["host"] == "{host}")
242 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
243 }
244
245 try:
246 data = await _run_queries(
247 influxdb_client,
248 org,
249 queries,
250 ts_keys=set(queries.keys()),
251 )
252 return MetricsResponse(success=True, message="Successfully retrieved CPU metrics", data=data)
253 except Exception as e:
254 logger.error(f"Error fetching CPU metrics: {e}")
255 return MetricsResponse(success=False, message=f"Error fetching CPU metrics: {e}")
256 finally:
257 await influxdb_client.close()
258
259
260 # ── Memory ───────────────────────────────────────────────────────────────
261
262
263 async def get_memory_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
264 """Retrieve memory metrics for a given host."""
265 try:
266 bucket = await _get_influxdb_bucket(session)
267 org = await get_influxdb_organization()
268 except Exception as e:
269 return MetricsResponse(success=False, message=str(e))
270
271 influxdb_client = await create_influxdb_client("InfluxDB")
272 rng = _range_clause(range_h)
273 win = _window(range_h)
274
275 queries = {
276 "mem_used": f"""from(bucket: "{bucket}") |> {rng}
277 |> filter(fn: (r) => r["_measurement"] == "mem")
278 |> filter(fn: (r) => r["_field"] == "used" or r["_field"] == "total")
279 |> filter(fn: (r) => r["host"] == "{host}")
280 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
281 "swap_total": f"""from(bucket: "{bucket}") |> {rng}
282 |> filter(fn: (r) => r["_measurement"] == "swap")
283 |> filter(fn: (r) => r["_field"] == "total")
284 |> filter(fn: (r) => r["host"] == "{host}")
285 |> last()""",
286 "swap_free": f"""from(bucket: "{bucket}") |> {rng}
287 |> filter(fn: (r) => r["_measurement"] == "swap")
288 |> filter(fn: (r) => r["_field"] == "free")
289 |> filter(fn: (r) => r["host"] == "{host}")
290 |> last()""",
291 }
292
293 try:
294 data = await _run_queries(influxdb_client, org, queries, ts_keys={"mem_used"})
295 return MetricsResponse(success=True, message="Successfully retrieved memory metrics", data=data)
296 except Exception as e:
297 logger.error(f"Error fetching memory metrics: {e}")
298 return MetricsResponse(success=False, message=f"Error fetching memory metrics: {e}")
299 finally:
300 await influxdb_client.close()
301
302
303 # ── Kernel ───────────────────────────────────────────────────────────────
304
305
306 async def get_kernel_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
307 """Retrieve kernel metrics for a given host."""
308 try:
309 bucket = await _get_influxdb_bucket(session)
310 org = await get_influxdb_organization()
311 except Exception as e:
312 return MetricsResponse(success=False, message=str(e))
313
314 influxdb_client = await create_influxdb_client("InfluxDB")
315 rng = _range_clause(range_h)
316 win = _window(range_h)
317
318 queries = {
319 "interrupts": f"""from(bucket: "{bucket}") |> {rng}
320 |> filter(fn: (r) => r["_measurement"] == "kernel")
321 |> filter(fn: (r) => r["_field"] == "interrupts")
322 |> filter(fn: (r) => r["host"] == "{host}")
323 |> derivative(unit: 1s, nonNegative: true)
324 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
325 "processes_forked": f"""from(bucket: "{bucket}") |> {rng}
326 |> filter(fn: (r) => r["_measurement"] == "kernel")
327 |> filter(fn: (r) => r["_field"] == "processes_forked")
328 |> filter(fn: (r) => r["host"] == "{host}")
329 |> derivative(unit: 1s, nonNegative: true)
330 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
331 }
332
333 try:
334 data = await _run_queries(
335 influxdb_client,
336 org,
337 queries,
338 ts_keys=set(queries.keys()),
339 )
340 return MetricsResponse(success=True, message="Successfully retrieved kernel metrics", data=data)
341 except Exception as e:
342 logger.error(f"Error fetching kernel metrics: {e}")
343 return MetricsResponse(success=False, message=f"Error fetching kernel metrics: {e}")
344 finally:
345 await influxdb_client.close()
346
347
348 # ── Disks ────────────────────────────────────────────────────────────────
349
350
351 async def get_disk_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
352 """Retrieve disk metrics for a given host."""
353 try:
354 bucket = await _get_influxdb_bucket(session)
355 org = await get_influxdb_organization()
356 except Exception as e:
357 return MetricsResponse(success=False, message=str(e))
358
359 influxdb_client = await create_influxdb_client("InfluxDB")
360 rng = _range_clause(range_h)
361 win = _window(range_h)
362
363 queries = {
364 "disk_total": f"""from(bucket: "{bucket}") |> {rng}
365 |> filter(fn: (r) => r["_measurement"] == "disk")
366 |> filter(fn: (r) => r["_field"] == "total")
367 |> filter(fn: (r) => r["host"] == "{host}")
368 |> last()""",
369 "disk_usage": f"""from(bucket: "{bucket}") |> {rng}
370 |> filter(fn: (r) => r["_measurement"] == "disk")
371 |> filter(fn: (r) => r["_field"] == "used_percent")
372 |> filter(fn: (r) => r["host"] == "{host}")
373 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
374 "disk_io": f"""from(bucket: "{bucket}") |> {rng}
375 |> filter(fn: (r) => r["_measurement"] == "diskio")
376 |> filter(fn: (r) => r["_field"] == "read_bytes" or r["_field"] == "write_bytes")
377 |> filter(fn: (r) => r["host"] == "{host}")
378 |> derivative(unit: 1s, nonNegative: true)
379 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
380 "inodes": f"""from(bucket: "{bucket}") |> {rng}
381 |> filter(fn: (r) => r["_measurement"] == "disk")
382 |> filter(fn: (r) => r["_field"] == "inodes_used" or r["_field"] == "inodes_total")
383 |> filter(fn: (r) => r["host"] == "{host}")
384 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
385 }
386
387 try:
388 query_api = influxdb_client.query_api()
389 data: dict[str, Any] = {}
390
391 for key, flux in queries.items():
392 result = await query_api.query(flux, org=org)
393 if key == "disk_total":
394 # Sum across all mount points
395 total = 0.0
396 for table in result:
397 for record in table.records:
398 try:
399 v = float(record.get_value())
400 if math.isfinite(v):
401 total += v
402 except (ValueError, TypeError):
403 pass
404 data[key] = total
405 elif key == "disk_io":
406 data[key] = _parse_ts_series(result, label_field="name")
407 elif key == "disk_usage" or key == "inodes":
408 data[key] = _parse_ts_series(result, label_field="path")
409 else:
410 data[key] = _parse_ts_series(result)
411
412 return MetricsResponse(success=True, message="Successfully retrieved disk metrics", data=data)
413 except Exception as e:
414 logger.error(f"Error fetching disk metrics: {e}")
415 return MetricsResponse(success=False, message=f"Error fetching disk metrics: {e}")
416 finally:
417 await influxdb_client.close()
418
419
420 # ── Processes ────────────────────────────────────────────────────────────
421
422
423 async def get_process_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
424 """Retrieve process metrics for a given host."""
425 try:
426 bucket = await _get_influxdb_bucket(session)
427 org = await get_influxdb_organization()
428 except Exception as e:
429 return MetricsResponse(success=False, message=str(e))
430
431 influxdb_client = await create_influxdb_client("InfluxDB")
432 rng = _range_clause(range_h)
433 win = _window(range_h)
434
435 queries = {
436 "status": f"""from(bucket: "{bucket}") |> {rng}
437 |> filter(fn: (r) => r["_measurement"] == "processes")
438 |> filter(fn: (r) => r["_field"] == "running" or r["_field"] == "sleeping" or r["_field"] == "zombies" or r["_field"] == "stopped" or r["_field"] == "blocked")
439 |> filter(fn: (r) => r["host"] == "{host}")
440 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
441 }
442
443 # Stat values
444 for field in ("running", "sleeping", "unknown", "zombies"):
445 queries[
446 field
447 ] = f"""from(bucket: "{bucket}") |> {rng}
448 |> filter(fn: (r) => r["_measurement"] == "processes")
449 |> filter(fn: (r) => r["_field"] == "{field}")
450 |> filter(fn: (r) => r["host"] == "{host}")
451 |> last()"""
452
453 try:
454 data = await _run_queries(influxdb_client, org, queries, ts_keys={"status"})
455 return MetricsResponse(success=True, message="Successfully retrieved process metrics", data=data)
456 except Exception as e:
457 logger.error(f"Error fetching process metrics: {e}")
458 return MetricsResponse(success=False, message=f"Error fetching process metrics: {e}")
459 finally:
460 await influxdb_client.close()
461
462
463 # ── Network ──────────────────────────────────────────────────────────────
464
465
466 async def get_network_metrics(host: str, range_h: str, session: AsyncSession) -> MetricsResponse:
467 """Retrieve network metrics for a given host."""
468 try:
469 bucket = await _get_influxdb_bucket(session)
470 org = await get_influxdb_organization()
471 except Exception as e:
472 return MetricsResponse(success=False, message=str(e))
473
474 influxdb_client = await create_influxdb_client("InfluxDB")
475 rng = _range_clause(range_h)
476 win = _window(range_h)
477
478 queries = {
479 "traffic": f"""from(bucket: "{bucket}") |> {rng}
480 |> filter(fn: (r) => r["_measurement"] == "net")
481 |> filter(fn: (r) => r["_field"] == "bytes_recv" or r["_field"] == "bytes_sent")
482 |> filter(fn: (r) => r["host"] == "{host}")
483 |> derivative(unit: 1s, nonNegative: true)
484 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
485 "tcp_established": f"""from(bucket: "{bucket}") |> {rng}
486 |> filter(fn: (r) => r["_measurement"] == "netstat")
487 |> filter(fn: (r) => r["_field"] == "tcp_established")
488 |> filter(fn: (r) => r["host"] == "{host}")
489 |> last()""",
490 "interface_errors": f"""from(bucket: "{bucket}") |> {rng}
491 |> filter(fn: (r) => r["_measurement"] == "net")
492 |> filter(fn: (r) => r["_field"] == "err_in" or r["_field"] == "err_out")
493 |> filter(fn: (r) => r["host"] == "{host}")
494 |> aggregateWindow(every: {win}, fn: mean, createEmpty: false)""",
495 }
496
497 try:
498 query_api = influxdb_client.query_api()
499 data: dict[str, Any] = {}
500
501 for key, flux in queries.items():
502 result = await query_api.query(flux, org=org)
503 if key == "tcp_established":
504 data[key] = _last_value(result)
505 else:
506 data[key] = _parse_ts_series(result, label_field="interface")
507
508 return MetricsResponse(success=True, message="Successfully retrieved network metrics", data=data)
509 except Exception as e:
510 logger.error(f"Error fetching network metrics: {e}")
511 return MetricsResponse(success=False, message=f"Error fetching network metrics: {e}")
512 finally:
513 await influxdb_client.close()