81 lines
2.5 KiB
Python
81 lines
2.5 KiB
Python
from __future__ import annotations
|
|
|
|
import logging
|
|
from typing import List
|
|
|
|
import requests
|
|
import urllib3
|
|
|
|
from .checker import CheckResult
|
|
|
|
logger = logging.getLogger("uptime_monitor.influx")
|
|
|
|
# The InfluxDB endpoint uses a self-signed certificate by design (per
|
|
# config, verify_ssl defaults to False for it) — silence the resulting
|
|
# urllib3 warning so it doesn't spam the logs on every write.
|
|
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
|
|
|
|
|
|
def _escape_tag(value) -> str:
|
|
return (
|
|
str(value)
|
|
.replace("\\", "\\\\")
|
|
.replace(",", "\\,")
|
|
.replace(" ", "\\ ")
|
|
.replace("=", "\\=")
|
|
)
|
|
|
|
|
|
def _escape_field_string(value) -> str:
|
|
return str(value).replace("\\", "\\\\").replace('"', '\\"')
|
|
|
|
|
|
def result_to_line(result: CheckResult, measurement: str = "uptime_check") -> str:
|
|
tags = (
|
|
f"endpoint={_escape_tag(result.endpoint_name)},"
|
|
f"host={_escape_tag(result.host)},"
|
|
f"scheme={_escape_tag(result.scheme)},"
|
|
f"port={result.port}"
|
|
)
|
|
fields = [
|
|
f"success={1 if result.success else 0}i",
|
|
f"expected_status={result.expected_status}i",
|
|
]
|
|
if result.status_code is not None:
|
|
fields.append(f"status_code={result.status_code}i")
|
|
if result.response_time_ms is not None:
|
|
fields.append(f"response_time_ms={result.response_time_ms}")
|
|
if result.resolved_ip:
|
|
fields.append(f'resolved_ip="{_escape_field_string(result.resolved_ip)}"')
|
|
if result.error:
|
|
fields.append(f'error="{_escape_field_string(result.error)}"')
|
|
|
|
return f"{measurement},{tags} {','.join(fields)} {result.timestamp_ns}"
|
|
|
|
|
|
def write_results(influx_cfg, results: List[CheckResult]) -> None:
|
|
if not results:
|
|
return
|
|
|
|
lines = "\n".join(result_to_line(r) for r in results)
|
|
url = f"{influx_cfg.url.rstrip('/')}/api/v2/write"
|
|
params = {"org": influx_cfg.org, "bucket": influx_cfg.bucket, "precision": "ns"}
|
|
headers = {
|
|
"Authorization": f"Token {influx_cfg.token}",
|
|
"Content-Type": "text/plain; charset=utf-8",
|
|
}
|
|
|
|
try:
|
|
resp = requests.post(
|
|
url,
|
|
params=params,
|
|
headers=headers,
|
|
data=lines.encode("utf-8"),
|
|
timeout=10,
|
|
verify=influx_cfg.verify_ssl,
|
|
)
|
|
resp.raise_for_status()
|
|
logger.info("wrote %d point(s) to influx bucket=%s", len(results), influx_cfg.bucket)
|
|
except requests.RequestException as e:
|
|
logger.error("failed to write to influx: %s", e)
|