723 lines
25 KiB
Python
723 lines
25 KiB
Python
"""
|
|
vw-eu-data-act-exporter
|
|
========================
|
|
Prometheus exporter for vehicle data from the VW Group EU Data Act Portal
|
|
(https://eu-data-act.drivesomethinggreater.com/).
|
|
|
|
For each configured VIN, a dedicated background thread logs into the portal,
|
|
downloads the latest available dataset, and exposes the fields it contains as
|
|
Prometheus metrics (label ``vin``) on a shared ``/metrics`` endpoint (the
|
|
classic Prometheus multi-target pattern via labels instead of separate
|
|
processes).
|
|
|
|
No raw data or ZIP files are written to disk - all state lives purely in
|
|
process memory and is lost on restart (health/last-scrape metrics then start
|
|
at 0 again). Optionally, the JSON extracted from every downloaded dataset -
|
|
including "_no_content_found.zip" placeholders, if they carry a JSON payload
|
|
- can be persisted to disk via ``config.persist_raw_json`` to help refine the
|
|
field mapping below. ZIP files themselves are never written to disk, only
|
|
the JSON extracted from them.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import io
|
|
import json
|
|
import logging
|
|
import random
|
|
import re
|
|
import sys
|
|
import threading
|
|
import time
|
|
import zipfile
|
|
from datetime import datetime
|
|
from pathlib import Path
|
|
from urllib.parse import urlencode, urljoin, urlparse
|
|
|
|
import requests
|
|
from bs4 import BeautifulSoup
|
|
from prometheus_client import Counter, Enum, Gauge, start_http_server
|
|
from prometheus_client.core import GaugeMetricFamily
|
|
from prometheus_client.registry import REGISTRY
|
|
|
|
import config
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Logging
|
|
# ---------------------------------------------------------------------------
|
|
logging.basicConfig(
|
|
level=getattr(logging, config.log_level.upper(), logging.INFO),
|
|
format="%(asctime)s %(levelname)-8s %(message)s",
|
|
datefmt="%H:%M:%S",
|
|
)
|
|
log = logging.getLogger("vw_eu_data_act_exporter")
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Portal & OIDC constants (brand "volkswagen" only)
|
|
# ---------------------------------------------------------------------------
|
|
BASE_URL = "https://eu-data-act.drivesomethinggreater.com"
|
|
IDENTITY_BASE = "https://identity.vwgroup.io"
|
|
OIDC_AUTHORIZE = IDENTITY_BASE + "/oidc/v1/authorize"
|
|
OIDC_SCOPE = "openid cars profile"
|
|
OIDC_REDIRECT = BASE_URL + "/login"
|
|
OIDC_CLIENT_ID = "9b58543e-1c15-4193-91d5-8a14145bebb0@apps_vw-dilab_com"
|
|
OIDC_STATE_KEY = "VOLKSWAGEN_PASSENGER_CARS"
|
|
DEFAULT_COUNTRY = "de"
|
|
DEFAULT_LANGUAGE = "en"
|
|
|
|
METADATA_PATH = "/proxy_api/euda-apim/datarequest/vehicles/{vin}/metadata/partial"
|
|
LIST_PATH = "/proxy_api/euda-apim/datadelivery/vehicles/{vin}/{identifier}/list"
|
|
DOWNLOAD_PATH = "/proxy_api/euda-apim/datadelivery/vehicles/{vin}/{identifier}/download"
|
|
|
|
NO_CONTENT_SUFFIX = "_no_content_found.zip"
|
|
HTTP_TIMEOUT = 30
|
|
|
|
IDENTIFIER_KEYS = {"identifier", "datarequestid", "id"}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# HTML / JS login-parsing helpers (adopted from vw_eu_data_act_downloader.py)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _extract_template_model(html: str) -> dict:
|
|
idx = html.find("templateModel")
|
|
if idx == -1:
|
|
return {}
|
|
brace = html.find("{", idx)
|
|
if brace == -1:
|
|
return {}
|
|
depth = 0
|
|
for i in range(brace, len(html)):
|
|
c = html[i]
|
|
if c == "{":
|
|
depth += 1
|
|
elif c == "}":
|
|
depth -= 1
|
|
if depth == 0:
|
|
try:
|
|
return json.loads(html[brace : i + 1])
|
|
except ValueError:
|
|
return {}
|
|
return {}
|
|
|
|
|
|
def _extract_csrf(html: str) -> str | None:
|
|
m = re.search(r"csrf_token\s*[:=]\s*['\"]([^'\"]+)['\"]", html)
|
|
return m.group(1) if m else None
|
|
|
|
|
|
def _parse_form_action(html: str, base_url: str) -> str | None:
|
|
soup = BeautifulSoup(html, "html.parser")
|
|
form = soup.find("form")
|
|
if not form:
|
|
return None
|
|
action = form.get("action")
|
|
if not action:
|
|
return None
|
|
return urljoin(base_url, action)
|
|
|
|
|
|
def _collect_login_fields(html: str) -> dict[str, str]:
|
|
soup = BeautifulSoup(html, "html.parser")
|
|
fields: dict[str, str] = {}
|
|
|
|
form = soup.find("form")
|
|
if form:
|
|
for inp in form.find_all("input"):
|
|
name = inp.get("name")
|
|
if name:
|
|
fields[name] = inp.get("value") or ""
|
|
|
|
model = _extract_template_model(html)
|
|
if model:
|
|
for key in ("hmac", "relayState"):
|
|
if model.get(key):
|
|
fields[key] = model[key]
|
|
email_val = (model.get("emailPasswordForm") or {}).get("email")
|
|
if email_val:
|
|
fields.setdefault("email", email_val)
|
|
|
|
csrf = _extract_csrf(html)
|
|
if csrf:
|
|
fields.setdefault("_csrf", csrf)
|
|
|
|
return fields
|
|
|
|
|
|
def _extract_login_error(html: str) -> str | None:
|
|
model = _extract_template_model(html)
|
|
err = model.get("error") or model.get("errorCode")
|
|
if isinstance(err, dict):
|
|
return err.get("text") or err.get("errorCode") or str(err)
|
|
return str(err) if err else None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Portal client
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _build_authorize_url() -> str:
|
|
state = f"{DEFAULT_COUNTRY}__{DEFAULT_LANGUAGE}__{OIDC_STATE_KEY}"
|
|
params = {
|
|
"client_id": OIDC_CLIENT_ID,
|
|
"response_type": "code",
|
|
"scope": OIDC_SCOPE,
|
|
"state": state,
|
|
"redirect_uri": OIDC_REDIRECT,
|
|
"prompt": "login",
|
|
}
|
|
return f"{OIDC_AUTHORIZE}?{urlencode(params)}"
|
|
|
|
|
|
def login(session: requests.Session, email: str, password: str) -> None:
|
|
"""Runs the full OIDC login flow. Raises RuntimeError on failure."""
|
|
|
|
try:
|
|
session.get(BASE_URL + "/", timeout=15)
|
|
except requests.RequestException as exc:
|
|
log.debug("Priming request failed (ignored): %s", exc)
|
|
|
|
authorize_url = _build_authorize_url()
|
|
resp = session.get(authorize_url, timeout=20)
|
|
resp.raise_for_status()
|
|
signin_url = resp.url
|
|
signin_html = resp.text
|
|
|
|
fields = _collect_login_fields(signin_html)
|
|
if "hmac" not in fields or "_csrf" not in fields:
|
|
raise RuntimeError(
|
|
f"Could not parse the sign-in form (fields found: {sorted(fields)})."
|
|
)
|
|
fields["email"] = email
|
|
action = _parse_form_action(signin_html, signin_url)
|
|
identifier_action = action or signin_url
|
|
resp = session.post(
|
|
identifier_action, data=fields, headers={"Referer": signin_url}, timeout=20
|
|
)
|
|
resp.raise_for_status()
|
|
authenticate_url = resp.url
|
|
authenticate_html = resp.text
|
|
|
|
fields2 = _collect_login_fields(authenticate_html)
|
|
if "hmac" not in fields2 or "_csrf" not in fields2:
|
|
err = _extract_login_error(authenticate_html)
|
|
raise RuntimeError(err or "IDP did not return a password form.")
|
|
fields2["email"] = email
|
|
fields2["password"] = password
|
|
|
|
action2 = _parse_form_action(authenticate_html, authenticate_url)
|
|
authenticate_action = action2 or authenticate_url.split("?", 1)[0]
|
|
resp = session.post(
|
|
authenticate_action,
|
|
data=fields2,
|
|
headers={"Referer": authenticate_url},
|
|
timeout=20,
|
|
)
|
|
|
|
if resp.status_code >= 400:
|
|
err = _extract_login_error(resp.text)
|
|
raise RuntimeError(err or f"Login rejected with HTTP {resp.status_code}.")
|
|
|
|
landing = resp.url
|
|
portal_host = urlparse(BASE_URL).netloc
|
|
if "signin-service" in landing or "/error" in landing:
|
|
err = _extract_login_error(resp.text)
|
|
raise RuntimeError(f"Login failed - {err or 'check email/password'}")
|
|
if urlparse(landing).netloc != portal_host:
|
|
raise RuntimeError(f"Login did not complete - landed at {landing!r}.")
|
|
|
|
|
|
def _api_get(session: requests.Session, url: str, *, headers: dict | None = None):
|
|
resp = session.get(url, headers=headers or {}, timeout=HTTP_TIMEOUT)
|
|
if resp.status_code >= 400:
|
|
raise RuntimeError(
|
|
f"API GET {url} -> HTTP {resp.status_code}: {resp.text[:300]}"
|
|
)
|
|
return resp.json()
|
|
|
|
|
|
def get_metadata(session: requests.Session, vin: str) -> dict:
|
|
url = BASE_URL + METADATA_PATH.format(vin=vin)
|
|
return _api_get(session, url)
|
|
|
|
|
|
def find_identifier(node) -> str | None:
|
|
if isinstance(node, dict):
|
|
for k, v in node.items():
|
|
if k.lower() in IDENTIFIER_KEYS and isinstance(v, str) and len(v) > 8:
|
|
return v
|
|
for v in node.values():
|
|
r = find_identifier(v)
|
|
if r:
|
|
return r
|
|
elif isinstance(node, list):
|
|
for v in node:
|
|
r = find_identifier(v)
|
|
if r:
|
|
return r
|
|
return None
|
|
|
|
|
|
def list_datasets(session: requests.Session, vin: str, identifier: str) -> list[dict]:
|
|
url = BASE_URL + LIST_PATH.format(vin=vin, identifier=identifier)
|
|
data = _api_get(session, url, headers={"type": "partial"})
|
|
return data if isinstance(data, list) else data.get("files", [])
|
|
|
|
|
|
def download_zip(
|
|
session: requests.Session, vin: str, identifier: str, name: str
|
|
) -> bytes:
|
|
url = BASE_URL + DOWNLOAD_PATH.format(vin=vin, identifier=identifier)
|
|
resp = session.get(url, headers={"filename": name, "type": "partial"}, timeout=60)
|
|
if resp.status_code >= 400:
|
|
raise RuntimeError(
|
|
f"Download {name} -> HTTP {resp.status_code}: {resp.text[:200]}"
|
|
)
|
|
return resp.content
|
|
|
|
|
|
def unzip_json(raw: bytes, name: str) -> dict:
|
|
with zipfile.ZipFile(io.BytesIO(raw)) as zf:
|
|
members = [m for m in zf.namelist() if m.lower().endswith(".json")]
|
|
if not members:
|
|
raise ValueError(f"No JSON file in {name}")
|
|
with zf.open(members[0]) as fh:
|
|
return json.loads(fh.read().decode("utf-8"))
|
|
|
|
|
|
def pick_latest_dataset(datasets: list[dict]) -> dict:
|
|
def sort_key(d: dict) -> str:
|
|
return d.get("createdOn") or d.get("created") or d.get("name") or ""
|
|
|
|
return sorted(datasets, key=sort_key)[-1]
|
|
|
|
|
|
def flatten_dataset(record: dict) -> dict[str, str]:
|
|
"""Builds a dataFieldName -> value mapping. On duplicates, the last
|
|
entry in the Data array wins."""
|
|
fields: dict[str, str] = {}
|
|
for entry in record.get("Data", []):
|
|
name = entry.get("dataFieldName")
|
|
if name is None:
|
|
continue
|
|
fields[name] = entry.get("value")
|
|
return fields
|
|
|
|
|
|
def persist_raw_dataset(vin: str, record: dict, *, suffix: str = "") -> None:
|
|
"""Writes the JSON extracted from a downloaded dataset to disk, so it
|
|
can be inspected later to refine the field mapping below. Only the
|
|
extracted JSON is ever written - never the ZIP it came from. Only
|
|
called when ``config.persist_raw_json`` is enabled."""
|
|
out_dir = Path(config.persist_raw_json_dir)
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
timestamp = datetime.utcnow().strftime("%Y%m%d%H%M%S")
|
|
tag = f"_{suffix}" if suffix else ""
|
|
dest = out_dir / f"{vin}{tag}_{timestamp}.json"
|
|
dest.write_text(json.dumps(record, indent=2, ensure_ascii=False), encoding="utf-8")
|
|
log.debug("Persisted raw dataset for VIN %s to %s", vin, dest)
|
|
|
|
|
|
def fetch_latest_dataset_fields(vin: str) -> dict[str, str] | None:
|
|
"""Logs in, downloads the latest dataset for *vin* and returns the
|
|
flattened fields. Returns ``None`` if no datasets are (yet) available -
|
|
this is explicitly NOT a failure. Raises an exception on login/HTTP/parse
|
|
errors (= failure)."""
|
|
|
|
session = requests.Session()
|
|
session.headers.update({"User-Agent": random.choice(config.user_agents)})
|
|
session.hooks["response"].append(
|
|
lambda resp, *args, **kwargs: c_http_requests_total.labels(vin=vin).inc()
|
|
)
|
|
|
|
login(session, config.vw_account["email"], config.vw_account["password"])
|
|
|
|
meta = get_metadata(session, vin)
|
|
identifier = find_identifier(meta)
|
|
if not identifier:
|
|
raise RuntimeError(f"No data-request identifier in metadata for VIN {vin}.")
|
|
|
|
all_datasets = list_datasets(session, vin, identifier)
|
|
no_content_datasets = [
|
|
d
|
|
for d in all_datasets
|
|
if (d.get("name") or d.get("fileName") or "").endswith(NO_CONTENT_SUFFIX)
|
|
]
|
|
datasets = [d for d in all_datasets if d not in no_content_datasets]
|
|
|
|
if config.persist_raw_json and no_content_datasets:
|
|
nc_latest = pick_latest_dataset(no_content_datasets)
|
|
nc_name = nc_latest.get("name") or nc_latest.get("fileName") or str(nc_latest)
|
|
try:
|
|
nc_raw = download_zip(session, vin, identifier, nc_name)
|
|
nc_record = unzip_json(nc_raw, nc_name)
|
|
except RuntimeError as exc:
|
|
log.warning(
|
|
"Could not download no-content placeholder for VIN %s: %s", vin, exc
|
|
)
|
|
except ValueError:
|
|
log.debug(
|
|
"No-content placeholder %s for VIN %s has no JSON payload - nothing to persist.",
|
|
nc_name,
|
|
vin,
|
|
)
|
|
else:
|
|
try:
|
|
persist_raw_dataset(vin, nc_record, suffix="no_content")
|
|
except OSError as exc:
|
|
log.warning(
|
|
"Could not persist no-content dataset for VIN %s: %s", vin, exc
|
|
)
|
|
|
|
if not datasets:
|
|
return None
|
|
|
|
latest = pick_latest_dataset(datasets)
|
|
name = latest.get("name") or latest.get("fileName") or str(latest)
|
|
raw = download_zip(session, vin, identifier, name)
|
|
record = unzip_json(raw, name)
|
|
|
|
if config.persist_raw_json:
|
|
try:
|
|
persist_raw_dataset(vin, record)
|
|
except OSError as exc:
|
|
log.warning("Could not persist raw dataset for VIN %s: %s", vin, exc)
|
|
|
|
return flatten_dataset(record)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Value conversion (VW API strings -> Prometheus types)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def parse_bool(raw) -> bool:
|
|
return str(raw).strip().lower() == "true"
|
|
|
|
|
|
def parse_float(raw) -> float:
|
|
return float(raw)
|
|
|
|
|
|
def parse_iso8601_to_unix(raw) -> float:
|
|
text = str(raw).strip()
|
|
if text.endswith("Z"):
|
|
text = text[:-1] + "+00:00"
|
|
return datetime.fromisoformat(text).timestamp()
|
|
|
|
|
|
def parse_bool_as_float(raw) -> float:
|
|
return 1.0 if parse_bool(raw) else 0.0
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Prometheus metrics
|
|
# ---------------------------------------------------------------------------
|
|
P = config.exporter_prefix
|
|
|
|
g_mileage_km = Gauge(f"{P}mileage_km", "Odometer reading", ["vin"])
|
|
g_hvsoc_percent = Gauge(
|
|
f"{P}hvsoc_percent", "HV battery state of charge in percent", ["vin"]
|
|
)
|
|
g_driver_present = Gauge(
|
|
f"{P}driver_present", "Driver detected in vehicle (1=yes)", ["vin"]
|
|
)
|
|
g_cruising_range_km = Gauge(f"{P}cruising_range_km", "Remaining range in km", ["vin"])
|
|
g_hvbattery_temp_max_celsius = Gauge(
|
|
f"{P}hvbattery_temperature_max_celsius",
|
|
"Max. HV battery temperature in Celsius",
|
|
["vin"],
|
|
)
|
|
g_hvbattery_temp_min_celsius = Gauge(
|
|
f"{P}hvbattery_temperature_min_celsius",
|
|
"Min. HV battery temperature in Celsius",
|
|
["vin"],
|
|
)
|
|
g_charge_power_kw = Gauge(
|
|
f"{P}charge_power_kw", "Current charging power in kW", ["vin"]
|
|
)
|
|
g_target_soc_percent = Gauge(
|
|
f"{P}target_soc_percent", "Charge target in percent", ["vin"]
|
|
)
|
|
g_position_longitude = Gauge(
|
|
f"{P}position_longitude", "Longitude of the last known position", ["vin"]
|
|
)
|
|
g_position_latitude = Gauge(
|
|
f"{P}position_latitude", "Latitude of the last known position", ["vin"]
|
|
)
|
|
g_position_created_timestamp_seconds = Gauge(
|
|
f"{P}position_created_timestamp_seconds",
|
|
"Unix timestamp of when the last position was recorded",
|
|
["vin"],
|
|
)
|
|
g_locked = Gauge(
|
|
f"{P}locked",
|
|
"Combined lock status of all doors, trunk and hood (1=locked)",
|
|
["vin"],
|
|
)
|
|
g_is_parked = Gauge(f"{P}is_parked", "Vehicle parked (1=yes)", ["vin"])
|
|
g_parking_brake_engaged = Gauge(
|
|
f"{P}parking_brake_engaged", "Parking brake engaged (1=yes)", ["vin"]
|
|
)
|
|
g_service_due_in_days = Gauge(
|
|
f"{P}service_due_in_days", "Remaining days until the next service", ["vin"]
|
|
)
|
|
g_last_vehicle_signal_timestamp_seconds = Gauge(
|
|
f"{P}last_vehicle_signal_timestamp_seconds",
|
|
"Unix timestamp of the last vehicle signal (carCapturedUTCTimestamp)",
|
|
["vin"],
|
|
)
|
|
|
|
g_last_successful_scrape_timestamp_seconds = Gauge(
|
|
f"{P}last_successful_scrape_timestamp_seconds",
|
|
"Unix timestamp of the last successful scrape with new data",
|
|
["vin"],
|
|
)
|
|
g_health = Gauge(
|
|
f"{P}health",
|
|
"1 once new data has been received at least once via the EU Data Act API, "
|
|
"0 initially and after every failure (not a failure: portal ok but no new data)",
|
|
["vin"],
|
|
)
|
|
c_http_requests_total = Counter(
|
|
f"{P}http_requests_total",
|
|
"Total number of HTTP requests made to the EU Data Act portal (login + API calls)",
|
|
["vin"],
|
|
)
|
|
|
|
# Enum metrics. The state space is best-effort (compiled from the sample
|
|
# dataset and comparable projects) plus an "UNKNOWN" fallback for unknown raw
|
|
# values - these are additionally logged as WARNING so the list can be
|
|
# extended if needed.
|
|
DRIVING_MODE_STATES = ["standard", "eco", "comfort", "sport", "individual", "UNKNOWN"]
|
|
CHARGING_STATE_STATES = [
|
|
"OFF",
|
|
"READY_FOR_CHARGING",
|
|
"NOT_READY_FOR_CHARGING",
|
|
"CHARGING",
|
|
"CONSERVING",
|
|
"ERROR",
|
|
"UNKNOWN",
|
|
]
|
|
PLUG_CONNECTION_STATES = ["CONNECTED", "DISCONNECTED", "UNKNOWN"]
|
|
SERVICE_TYPE_STATES = [
|
|
"SERVICE_TYPE_INSPECTION",
|
|
"SERVICE_TYPE_OIL_CHANGE",
|
|
"SERVICE_TYPE_BRAKE_FLUID",
|
|
"SERVICE_TYPE_TIMING_BELT",
|
|
"UNKNOWN",
|
|
]
|
|
|
|
e_driving_mode = Enum(
|
|
f"{P}driving_mode", "Active driving mode", ["vin"], states=DRIVING_MODE_STATES
|
|
)
|
|
e_charging_state = Enum(
|
|
f"{P}charging_state",
|
|
"Current charging state",
|
|
["vin"],
|
|
states=CHARGING_STATE_STATES,
|
|
)
|
|
e_plug_connection_state = Enum(
|
|
f"{P}plug_connection_state",
|
|
"Plug connection status",
|
|
["vin"],
|
|
states=PLUG_CONNECTION_STATES,
|
|
)
|
|
e_service_type = Enum(
|
|
f"{P}next_service_type",
|
|
"Next due service type",
|
|
["vin"],
|
|
states=SERVICE_TYPE_STATES,
|
|
)
|
|
|
|
# Field name (dataFieldName in the dataset) -> (Gauge, conversion function)
|
|
GAUGE_FIELD_MAP: list[tuple[str, Gauge, callable]] = [
|
|
("mileage_info.value", g_mileage_km, parse_float),
|
|
("hvsoc_info.value", g_hvsoc_percent, parse_float),
|
|
("Driver Presence", g_driver_present, parse_bool_as_float),
|
|
("batteryStatus.cruisingRange.range", g_cruising_range_km, parse_float),
|
|
(
|
|
"hvbatterytemperature_info.max_temperature.value",
|
|
g_hvbattery_temp_max_celsius,
|
|
parse_float,
|
|
),
|
|
(
|
|
"hvbatterytemperature_info.min_temperature.value",
|
|
g_hvbattery_temp_min_celsius,
|
|
parse_float,
|
|
),
|
|
("chargingStatus.chargePower_kW", g_charge_power_kw, parse_float),
|
|
("targetSoc_pct", g_target_soc_percent, parse_float),
|
|
("positionCreated", g_position_created_timestamp_seconds, parse_iso8601_to_unix),
|
|
("isParked", g_is_parked, parse_bool_as_float),
|
|
("parking_brake_info.value", g_parking_brake_engaged, parse_bool_as_float),
|
|
("service_maintenance_info.due_in_time.value", g_service_due_in_days, parse_float),
|
|
(
|
|
"carCapturedUTCTimestamp",
|
|
g_last_vehicle_signal_timestamp_seconds,
|
|
parse_iso8601_to_unix,
|
|
),
|
|
]
|
|
|
|
ENUM_FIELD_MAP: list[tuple[str, Enum, list[str]]] = [
|
|
("drivingMode", e_driving_mode, DRIVING_MODE_STATES),
|
|
("chargingStatus.currentChargeState", e_charging_state, CHARGING_STATE_STATES),
|
|
(
|
|
"plugStatusItem.plugConnectionState",
|
|
e_plug_connection_state,
|
|
PLUG_CONNECTION_STATES,
|
|
),
|
|
("service_maintenance_info.service_type", e_service_type, SERVICE_TYPE_STATES),
|
|
]
|
|
|
|
# Lock components -> internal name. If a component is missing from a dataset
|
|
# entirely (even across all scrapes so far), it is treated as "LOCKED" (safe
|
|
# default), so a single door actively reporting locked does not on its own
|
|
# report "vehicle locked" if another component is still unknown.
|
|
LOCK_FIELDS = {
|
|
"door_info.front_left.door_lock_status.value": "front_left",
|
|
"door_info.front_right.door_lock_status.value": "front_right",
|
|
"door_info.rear_left.door_lock_status.value": "rear_left",
|
|
"door_info.rear_right.door_lock_status.value": "rear_right",
|
|
"trunk_lid_info.trunk_lid_lock_status.value": "trunk",
|
|
"hood_info.hood_lock_status.value": "hood",
|
|
}
|
|
|
|
|
|
class UptimeCollector:
|
|
"""Computes uptime freshly on every scrape instead of maintaining a
|
|
ticker thread."""
|
|
|
|
def __init__(self, start_time: float) -> None:
|
|
self._start_time = start_time
|
|
|
|
def collect(self):
|
|
g = GaugeMetricFamily(
|
|
f"{P}uptime_seconds", "Uptime of the exporter process in seconds"
|
|
)
|
|
g.add_metric([], time.time() - self._start_time)
|
|
yield g
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Per-VIN background poller
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class VehiclePoller:
|
|
def __init__(self, vin: str) -> None:
|
|
self.vin = vin
|
|
self.lock_components: dict[str, str] = {
|
|
name: "LOCKED" for name in LOCK_FIELDS.values()
|
|
}
|
|
self.last_signal_ts_raw: str | None = None
|
|
|
|
def apply_fields(self, fields: dict[str, str]) -> None:
|
|
vin = self.vin
|
|
|
|
for key, gauge, transform in GAUGE_FIELD_MAP:
|
|
if key not in fields:
|
|
continue
|
|
try:
|
|
gauge.labels(vin=vin).set(transform(fields[key]))
|
|
except (TypeError, ValueError) as exc:
|
|
log.warning("Could not parse field %s (VIN %s): %s", key, vin, exc)
|
|
|
|
for key, enum_metric, states in ENUM_FIELD_MAP:
|
|
if key not in fields:
|
|
continue
|
|
raw = str(fields[key])
|
|
state = raw if raw in states else "UNKNOWN"
|
|
if state == "UNKNOWN":
|
|
log.warning("Unknown enum value for %s (VIN %s): %r", key, vin, raw)
|
|
enum_metric.labels(vin=vin).state(state)
|
|
|
|
for key, comp_name in LOCK_FIELDS.items():
|
|
if key in fields:
|
|
self.lock_components[comp_name] = str(fields[key])
|
|
locked = (
|
|
0.0 if any(v == "UNLOCKED" for v in self.lock_components.values()) else 1.0
|
|
)
|
|
g_locked.labels(vin=vin).set(locked)
|
|
|
|
if "longitude" in fields and "latitude" in fields:
|
|
try:
|
|
g_position_longitude.labels(vin=vin).set(
|
|
parse_float(fields["longitude"])
|
|
)
|
|
g_position_latitude.labels(vin=vin).set(parse_float(fields["latitude"]))
|
|
except (TypeError, ValueError) as exc:
|
|
log.warning("Could not parse position (VIN %s): %s", vin, exc)
|
|
|
|
new_ts = fields.get("carCapturedUTCTimestamp")
|
|
if new_ts is not None and new_ts != self.last_signal_ts_raw:
|
|
self.last_signal_ts_raw = new_ts
|
|
g_last_successful_scrape_timestamp_seconds.labels(vin=vin).set(time.time())
|
|
g_health.labels(vin=vin).set(1.0)
|
|
|
|
def run_once(self) -> None:
|
|
fields = fetch_latest_dataset_fields(self.vin)
|
|
if fields is None:
|
|
log.info("No new datasets available for VIN %s.", self.vin)
|
|
return
|
|
self.apply_fields(fields)
|
|
|
|
def loop(self, stop_event: threading.Event) -> None:
|
|
g_health.labels(vin=self.vin).set(0.0)
|
|
while not stop_event.is_set():
|
|
try:
|
|
self.run_once()
|
|
except Exception as exc: # login/HTTP/parse errors = failure
|
|
log.error("Scrape for VIN %s failed: %s", self.vin, exc)
|
|
g_health.labels(vin=self.vin).set(0.0)
|
|
stop_event.wait(config.scrape_interval_minutes * 60)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Main
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def main() -> int:
|
|
if not config.vins:
|
|
log.error("config.vins is empty - no vehicles configured to scrape.")
|
|
return 1
|
|
|
|
REGISTRY.register(UptimeCollector(time.time()))
|
|
|
|
stop_event = threading.Event()
|
|
pollers = [VehiclePoller(vin) for vin in config.vins]
|
|
threads = []
|
|
for poller in pollers:
|
|
t = threading.Thread(
|
|
target=poller.loop,
|
|
args=(stop_event,),
|
|
daemon=True,
|
|
name=f"poll-{poller.vin}",
|
|
)
|
|
t.start()
|
|
threads.append(t)
|
|
|
|
start_http_server(config.serverPort, addr=config.hostName)
|
|
log.info(
|
|
"vw-eu-data-act-exporter running on http://%s:%s/metrics (%d VIN(s), interval %d min)",
|
|
config.hostName,
|
|
config.serverPort,
|
|
len(config.vins),
|
|
config.scrape_interval_minutes,
|
|
)
|
|
|
|
try:
|
|
while True:
|
|
time.sleep(3600)
|
|
except KeyboardInterrupt:
|
|
log.info("Shutting down...")
|
|
stop_event.set()
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|