| 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() |