Author SHA1 Message Date
jensandClaude Sonnet 4.6 6fb9e624e4 collect: fix stale vehicle ref after token re-login; increase update verbosity
After a full re-login (doLogin), the library rebuilds its vehicle list with
new objects. The old vehicle reference in the collect loop was never updated
again, causing data to silently freeze. Re-select the vehicle from wc.vehicles
after every update() call to stay on the current object.

Also: update() now returns and logs its bool result; the record log line shows
which top-level keys changed each cycle (or "none" when the backend returned
identical data).

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-30 12:31:30 +02:00
jensandClaude Sonnet 4.6 b75c4fe411 server: guard against None snapshots crashing the accept loop
If extract_all() returns None (Dashboard init failure), appending it to
store["records"] would cause _accept_loop to crash with AttributeError on
the next client connection.  Since only OSError was caught, the thread
died permanently — no further clients could join push_clients, so live
broadcasts reached nobody.

- collect.py: skip None snapshots instead of appending them
- network.py: filter None records in _accept_loop; catch all exceptions
  during history send (not just OSError) so the thread survives
- storage_helpers.py: filter None records when loading from disk

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-30 00:42:31 +02:00
jensandClaude Sonnet 4.6 dfd5ef7f56 server: remove diff switch, always send diffs
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 22:09:51 +02:00
jensandClaude Sonnet 4.6 31d9ce2311 data_model: fix apply_procedural paths for extract_all snapshot format
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 21:26:47 +02:00
jensandClaude Sonnet 4.6 1c02869bd3 gui_client: keep collapsed nodes collapsed on record updates
collect: capture apply_procedural return value

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 21:23:28 +02:00
jensandClaude Sonnet 4.6 73f6b85c05 gui_client: highlight changed rows on each new record
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 21:15:59 +02:00
jensandClaude Sonnet 4.6 74379e55d6 gui_client: add "Human Readable State" column populated from str_state
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 21:03:54 +02:00
jensandClaude Sonnet 4.6 78b6108658 gui_client: remove diff mode switch, always use diff mode
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 21:01:16 +02:00
jensandClaude Sonnet 4.6 6e340a9582 data_model: add extract_all() using Dashboard instruments as primary snapshot
Replaces the domain-based collect_snapshot with extract_all(), which iterates
all supported instruments via vw_dashboard.Dashboard. Each entry is stored as
{value, unit, str_state} grouped by component type. Units are validated against
a whitelist and normalised (d→days, °C→degC); unrecognised units are dropped.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 20:55:07 +02:00
jensandClaude Sonnet 4.6 f9dd96269e collect: revert last_broadcast to snapshot, not payload
Diffs should always be computed against the last full snapshot,
not the last broadcast payload.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 11:57:15 +02:00
jensandClaude Sonnet 4.6 561c610b92 collect: fix last_broadcast tracking and diff condition
Always compute the payload (diff or full snapshot) regardless of
last_broadcast state, and store the broadcast payload — not the raw
snapshot — as last_broadcast so subsequent diffs are based on what
was actually sent.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 11:46:54 +02:00
jensandClaude Sonnet 4.6 df944cdefb install_server: install all deps via requirements.txt, patch after
Reverts the split-install approach. volkswagencarnet stays in
server/requirements.txt; the patch is applied in a single step
after pip install completes.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 09:47:12 +02:00
jensandClaude Sonnet 4.6 4e2a22d35e installers: move to project root, share a single .venv
client/install.sh → install_client.sh
server/install.sh → install_server.sh

Both scripts now resolve REPO_DIR as their own directory (repo root)
and use a shared .venv at the repo root instead of per-subdirectory
virtual environments.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 09:30:07 +02:00
jensandClaude Sonnet 4.6 8ff9463e76 server: install volkswagencarnet from GitHub + apply auth patch
requirements.txt now references the upstream GitHub repo pinned to
commit c30dc37 (the last upstream commit before our fix), so the source
is portable and not machine-specific.

install.sh applies cariad-hybrid-auth-fix.patch to the installed
site-packages after pip runs, using `patch -N` so it is idempotent on
repeat installs.  The patch file is checked in to the repo so install.sh
can always find it.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-28 23:03:44 +02:00
11 changed files with 509 additions and 111 deletions
+283
View File
@@ -0,0 +1,283 @@
diff --git a/volkswagencarnet/vw_connection.py b/volkswagencarnet/vw_connection.py
index aac7fcc..944e8e1 100644
--- a/volkswagencarnet/vw_connection.py
+++ b/volkswagencarnet/vw_connection.py
@@ -4,10 +4,12 @@
from __future__ import annotations
import asyncio
+import base64
from datetime import UTC, datetime, timedelta
import hashlib
import logging
from random import randint, random
+import secrets
from urllib.parse import parse_qs, urljoin, urlparse
from typing import Dict, Optional
@@ -138,26 +140,43 @@ class Connection:
)
return await req.json()
- async def get_authorization_page(self, authorization_endpoint: str) -> str:
- """Get authorization page (login page)."""
- # https://identity.vwgroup.io/oidc/v1/authorize?nonce={NONCE}&state={STATE}&response_type={TOKEN_TYPES}&scope={SCOPE}&redirect_uri={APP_URI}&client_id={CLIENT_ID}
- # https://identity.vwgroup.io/oidc/v1/authorize?client_id={CLIENT_ID}&scope={SCOPE}&response_type={TOKEN_TYPES}&redirect_uri={APP_URI}
+ @staticmethod
+ def _generate_pkce_pair() -> tuple[str, str]:
+ """Return (code_verifier, code_challenge) for PKCE S256."""
+ verifier = base64.urlsafe_b64encode(secrets.token_bytes(32)).rstrip(b"=").decode()
+ digest = hashlib.sha256(verifier.encode()).digest()
+ challenge = base64.urlsafe_b64encode(digest).rstrip(b"=").decode()
+ return verifier, challenge
+
+ async def get_authorization_page(
+ self, authorization_endpoint: str, extra_params: dict | None = None
+ ) -> tuple[str, str]:
+ """Get authorization page (login page).
+
+ Returns:
+ (page_html, page_url) — the URL is needed as the POST target when
+ the form has no explicit action attribute.
+ """
_LOGGER.debug('Requesting authorization page from "%s"', authorization_endpoint)
self._session_auth_headers.pop("Referer", None)
self._session_auth_headers.pop("Origin", None)
_LOGGER.debug('Request headers: "%s"', self._session_auth_headers)
+ params = {
+ "redirect_uri": APP_URI,
+ "response_type": CLIENT_TOKEN_TYPES,
+ "client_id": CLIENT_ID,
+ "scope": CLIENT_SCOPE,
+ }
+ if extra_params:
+ params.update(extra_params)
+
try:
req = await self._session.get(
url=authorization_endpoint,
headers=self._session_auth_headers,
allow_redirects=False,
- params={
- "redirect_uri": APP_URI,
- "response_type": CLIENT_TOKEN_TYPES,
- "client_id": CLIENT_ID,
- "scope": CLIENT_SCOPE,
- },
+ params=params,
)
# Check if the response contains a redirect location
@@ -177,15 +196,20 @@ class Connection:
_LOGGER.info("Authorization error: %s", error_description)
raise AuthenticationError(f"{error_msg}: {error_description}")
- # If redirected, fetch the new location
- req = await self._session.get(
- url=ref, headers=self._session_auth_headers, allow_redirects=False
- )
-
- if req.status != 200:
- raise AuthenticationError("Failed to fetch authorization endpoint")
-
- return await req.text()
+ # Follow redirects to the login page (may be more than one hop)
+ while True:
+ req = await self._session.get(
+ url=ref, headers=self._session_auth_headers, allow_redirects=False
+ )
+ if req.status == 200:
+ return await req.text(), ref
+ if req.status in (301, 302, 303, 307, 308):
+ next_loc = req.headers.get("Location")
+ if not next_loc:
+ raise AuthenticationError("Missing Location header during login redirect chain")
+ ref = urljoin(ref, next_loc)
+ continue
+ raise AuthenticationError(f"Failed to fetch authorization endpoint: HTTP {req.status}")
except Exception as e:
_LOGGER.warning("Error during fetching authorization page: %s", str(e))
@@ -200,6 +224,26 @@ class Connection:
return None
return state_input["value"]
+ def extract_login_form_data(self, page_content: str) -> dict:
+ """Extract all hidden fields and submit button values from the login form."""
+ soup = BeautifulSoup(page_content, "html.parser")
+ form = soup.select_one("form")
+ if not form:
+ return {}
+ form_data = {}
+ for inp in form.find_all("input"):
+ name = inp.get("name")
+ value = inp.get("value", "")
+ if name and inp.get("type") in ("hidden", None):
+ form_data[name] = value
+ # Include submit button value if it carries an action name (e.g. name="action" value="default")
+ for btn in form.find_all("button", type="submit"):
+ name = btn.get("name")
+ value = btn.get("value", "")
+ if name:
+ form_data[name] = value
+ return form_data
+
async def post_form(
self, session, url: str, headers: dict, form_data: dict, redirect: bool = True
) -> str:
@@ -227,7 +271,8 @@ class Connection:
if error_code == "wrong-email-credentials":
raise AuthenticationError("Wrong username or password")
- # Unknown 400 error
+ # Unknown 400 error — log the response body to aid debugging
+ _LOGGER.debug("Login 400 response body: %s", page_content[:2000])
raise AuthenticationError(
"Login form validation failed with unknown 400 error"
)
@@ -303,10 +348,22 @@ class Connection:
authorization_endpoint = openid_config["authorization_endpoint"]
auth_issuer = openid_config["issuer"]
- # Get authorization page
- authorization_page = await self.get_authorization_page(authorization_endpoint)
+ # PKCE + nonce (nonce required when id_token appears in response_type)
+ code_verifier, code_challenge = self._generate_pkce_pair()
+ nonce = base64.urlsafe_b64encode(secrets.token_bytes(16)).rstrip(b"=").decode()
+
+ # Get authorization page; returns (html, page_url) — use page_url as POST target
+ # because the form has no action attribute and implicitly submits to itself.
+ authorization_page, login_url = await self.get_authorization_page(
+ authorization_endpoint,
+ extra_params={
+ "code_challenge": code_challenge,
+ "code_challenge_method": "S256",
+ "nonce": nonce,
+ },
+ )
- # Extract form data
+ # Extract state token (still needed for the login URL)
state_token = self.extract_state_token(authorization_page)
if not state_token:
@@ -316,13 +373,12 @@ class Connection:
)
raise AuthenticationError("Invalid login page - missing state token")
- # Do login
- login_form = {
- "username": self._session_auth_username,
- "password": self._session_auth_password,
- "state": state_token,
- }
- login_url = f"{auth_issuer}/u/login?state={state_token}"
+ # Build form from all hidden fields on the page so Auth0 gets every
+ # required field (e.g. csrf tokens, action), then overlay credentials.
+ login_form = self.extract_login_form_data(authorization_page)
+ login_form["username"] = self._session_auth_username
+ login_form["password"] = self._session_auth_password
+ login_form.setdefault("state", state_token)
redirect_location = await self.post_form(
self._session,
@@ -337,36 +393,40 @@ class Connection:
self._session, auth_issuer, redirect_location
)
- jwt_auth_code = parse_qs(urlparse(redirect_response).query)["code"][0]
- return jwt_auth_code
+ # Parse both query string and fragment — hybrid flow puts tokens in the fragment
+ all_params = {
+ **parse_qs(urlparse(redirect_response).query),
+ **parse_qs(urlparse(redirect_response).fragment),
+ }
+ jwt_auth_code = all_params["code"][0]
+ callback_id_token = all_params.get("id_token", [None])[0]
+ callback_access_token = all_params.get("access_token", [None])[0]
+ return jwt_auth_code, code_verifier, callback_id_token, callback_access_token
async def _exchange_code_for_tokens(
- self, auth_code: str, token_endpoint: str
+ self, auth_code: str, token_endpoint: str, code_verifier: str | None = None
) -> dict:
- """Exchange authorization code for access tokens.
-
- Args:
- auth_code: Authorization code from login flow
- token_endpoint: Token endpoint URL
-
- Returns:
- Dictionary containing tokens
-
- Raises:
- AuthenticationError: If token exchange fails
- """
+ """Exchange authorization code for access tokens via the CARIAD BFF token endpoint."""
token_body = {
"client_id": CLIENT_ID,
"grant_type": "authorization_code",
"code": auth_code,
"redirect_uri": APP_URI,
}
+ if code_verifier:
+ token_body["code_verifier"] = code_verifier
+
+ token_headers = {
+ "Accept-Encoding": "gzip, deflate, br",
+ "Connection": "keep-alive",
+ "Content-Type": "application/x-www-form-urlencoded",
+ "User-Agent": USER_AGENT,
+ "x-android-package-name": ANDROID_PACKAGE_NAME,
+ }
- # Token endpoint
token_response = await self.post_form(
- self._session, token_endpoint, self._session_auth_headers, token_body
+ self._session, token_endpoint, token_headers, token_body
)
-
return json_loads(token_response)
async def _login(self) -> bool:
@@ -385,11 +445,21 @@ class Connection:
openid_config = await self.get_openid_config()
token_endpoint = openid_config["token_endpoint"]
- # Get authorization code
- auth_code = await self._get_authorization_code(openid_config)
-
- # Exchange code for tokens
- tokens = await self._exchange_code_for_tokens(auth_code, token_endpoint)
+ # Get authorization code; hybrid response_type also delivers tokens directly
+ auth_code, code_verifier, cb_id_token, cb_access_token = await self._get_authorization_code(openid_config)
+
+ # Hybrid flow: tokens arrive in the callback URL — use them directly.
+ # The CARIAD BFF token endpoint requires a client_secret embedded in the
+ # official VW app, so we skip it when the callback already carries the tokens.
+ if cb_id_token and cb_access_token:
+ _LOGGER.debug("Using tokens from hybrid flow callback")
+ tokens = {
+ "access_token": cb_access_token,
+ "id_token": cb_id_token,
+ "token_type": "Bearer",
+ }
+ else:
+ tokens = await self._exchange_code_for_tokens(auth_code, token_endpoint, code_verifier)
# Validate token structure
required_keys = ["access_token", "id_token", "token_type"]
diff --git a/volkswagencarnet/vw_const.py b/volkswagencarnet/vw_const.py
index 17666a1..711fcf9 100644
--- a/volkswagencarnet/vw_const.py
+++ b/volkswagencarnet/vw_const.py
@@ -7,7 +7,7 @@ COUNTRY = "DE"
# Data used in communication
CLIENT_ID = "a24fba63-34b3-4d43-b181-942111e6bda8@apps_vw-dilab_com"
CLIENT_SCOPE = "openid profile badge cars dealers vin"
-CLIENT_TOKEN_TYPES = "code"
+CLIENT_TOKEN_TYPES = "code id_token token"
USER_AGENT = "Volkswagen/3.61.0-android/14"
APP_URI = "weconnect://authenticated"
+70 -50
View File
@@ -31,7 +31,7 @@ from pathlib import Path
from PyQt5.QtCore import QObject, QTimer, Qt, pyqtSignal from PyQt5.QtCore import QObject, QTimer, Qt, pyqtSignal
from PyQt5.QtGui import QColor, QFont from PyQt5.QtGui import QColor, QFont
from PyQt5.QtWidgets import ( from PyQt5.QtWidgets import (
QApplication, QCheckBox, QComboBox, QFormLayout, QGridLayout, QGroupBox, QApplication, QComboBox, QFormLayout, QGridLayout, QGroupBox,
QHBoxLayout, QLabel, QLineEdit, QMainWindow, QPushButton, QScrollArea, QHBoxLayout, QLabel, QLineEdit, QMainWindow, QPushButton, QScrollArea,
QSizePolicy, QSpinBox, QStatusBar, QTabWidget, QTreeWidget, QTreeWidgetItem, QSizePolicy, QSpinBox, QStatusBar, QTabWidget, QTreeWidget, QTreeWidgetItem,
QVBoxLayout, QWidget, QVBoxLayout, QWidget,
@@ -124,20 +124,18 @@ class TcpReader(QObject):
def __init__(self): def __init__(self):
super().__init__() super().__init__()
self._sock = None self._sock = None
self._running = False self._running = False
self._diff_mode = False
self._state: dict = {} self._state: dict = {}
def connect_to(self, host: str, port: int, diff_mode: bool = False) -> str | None: def connect_to(self, host: str, port: int) -> str | None:
"""Open connection; returns error string or None on success.""" """Open connection; returns error string or None on success."""
try: try:
self._sock = socket.create_connection((host, port), timeout=5) self._sock = socket.create_connection((host, port), timeout=5)
self._sock.settimeout(None) self._sock.settimeout(None)
except OSError as exc: except OSError as exc:
return str(exc) return str(exc)
self._diff_mode = diff_mode self._state = {}
self._state = {}
self._running = True self._running = True
threading.Thread(target=self._read_loop, daemon=True).start() threading.Thread(target=self._read_loop, daemon=True).start()
return None return None
@@ -165,11 +163,8 @@ class TcpReader(QObject):
if line: if line:
try: try:
data = json.loads(line) data = json.loads(line)
if self._diff_mode: self._state = jay_merge_full(self._state, data)
self._state = jay_merge_full(self._state, data) self.message.emit(copy.deepcopy(self._state))
self.message.emit(copy.deepcopy(self._state))
else:
self.message.emit(data)
except json.JSONDecodeError: except json.JSONDecodeError:
pass pass
except OSError: except OSError:
@@ -185,6 +180,23 @@ def _is_phys(v) -> bool:
return isinstance(v, dict) and "value" in v and "unit" in v return isinstance(v, dict) and "value" in v and "unit" in v
_CHANGED_BG = QColor("#1e3a00")
def _flatten_values(data: dict, prefix: str = "") -> dict[str, object]:
"""Flat dict of path → value for every physical leaf, using '/' separator."""
result = {}
for k, v in data.items():
if k == "ts":
continue
path = f"{prefix}/{k}" if prefix else k
if _is_phys(v):
result[path] = v.get("value")
elif isinstance(v, dict):
result.update(_flatten_values(v, path))
return result
def _extract_phys(data: dict, prefix: str = "") -> dict[str, dict]: def _extract_phys(data: dict, prefix: str = "") -> dict[str, dict]:
"""Return a flat dict of path → {value, unit} for every physical leaf.""" """Return a flat dict of path → {value, unit} for every physical leaf."""
result = {} result = {}
@@ -217,11 +229,9 @@ class ConnectorTab(QWidget):
self._port = QSpinBox() self._port = QSpinBox()
self._port.setRange(1, 65535) self._port.setRange(1, 65535)
self._port.setValue(9999) self._port.setValue(9999)
self._diff_cb = QCheckBox("Diff mode")
self._load_conn_settings() self._load_conn_settings()
form.addRow("Host:", self._host) form.addRow("Host:", self._host)
form.addRow("Port:", self._port) form.addRow("Port:", self._port)
form.addRow(self._diff_cb)
btn_row = QHBoxLayout() btn_row = QHBoxLayout()
self._btn = QPushButton("Connect") self._btn = QPushButton("Connect")
@@ -244,10 +254,6 @@ class ConnectorTab(QWidget):
else: else:
self.disconnect_requested.emit() self.disconnect_requested.emit()
@property
def diff_mode(self) -> bool:
return self._diff_cb.isChecked()
def _save_conn_settings(self): def _save_conn_settings(self):
try: try:
_SETTINGS_PATH.parent.mkdir(parents=True, exist_ok=True) _SETTINGS_PATH.parent.mkdir(parents=True, exist_ok=True)
@@ -256,9 +262,8 @@ class ConnectorTab(QWidget):
data = json.loads(_SETTINGS_PATH.read_text()) data = json.loads(_SETTINGS_PATH.read_text())
except (OSError, json.JSONDecodeError): except (OSError, json.JSONDecodeError):
pass pass
data["host"] = self._host.text() data["host"] = self._host.text()
data["port"] = self._port.value() data["port"] = self._port.value()
data["diff_mode"] = self._diff_cb.isChecked()
_SETTINGS_PATH.write_text(json.dumps(data, indent=2)) _SETTINGS_PATH.write_text(json.dumps(data, indent=2))
except OSError: except OSError:
pass pass
@@ -270,8 +275,6 @@ class ConnectorTab(QWidget):
self._host.setText(data["host"]) self._host.setText(data["host"])
if "port" in data: if "port" in data:
self._port.setValue(int(data["port"])) self._port.setValue(int(data["port"]))
if "diff_mode" in data:
self._diff_cb.setChecked(bool(data["diff_mode"]))
except (OSError, json.JSONDecodeError, KeyError, ValueError): except (OSError, json.JSONDecodeError, KeyError, ValueError):
pass pass
@@ -307,19 +310,25 @@ class DashboardTab(QWidget):
layout.addWidget(self._ts) layout.addWidget(self._ts)
self._tree = QTreeWidget() self._tree = QTreeWidget()
self._tree.setColumnCount(3) self._tree.setColumnCount(4)
self._tree.setHeaderLabels(["Field", "Value", "Unit"]) self._tree.setHeaderLabels(["Field", "Value", "Unit", "Human Readable State"])
self._tree.setColumnWidth(0, 320) self._tree.setColumnWidth(0, 320)
self._tree.setColumnWidth(1, 180) self._tree.setColumnWidth(1, 120)
self._tree.setColumnWidth(2, 80) self._tree.setColumnWidth(2, 60)
self._tree.setColumnWidth(3, 200)
self._tree.setAlternatingRowColors(True) self._tree.setAlternatingRowColors(True)
self._tree.setRootIsDecorated(True) self._tree.setRootIsDecorated(True)
layout.addWidget(self._tree) layout.addWidget(self._tree)
self._prev_flat: dict = {}
def update(self, snapshot: dict): def update(self, snapshot: dict):
self._ts.setText(f"Last update: {snapshot.get('ts', '')}") self._ts.setText(f"Last update: {snapshot.get('ts', '')}")
expanded = self._expanded_paths() curr_flat = _flatten_values(snapshot)
changed = {p for p, v in curr_flat.items() if v != self._prev_flat.get(p)}
self._prev_flat = curr_flat
collapsed = self._collapsed_paths()
self._tree.clear() self._tree.clear()
bold = QFont() bold = QFont()
@@ -333,67 +342,78 @@ class DashboardTab(QWidget):
d_item.setFont(0, bold) d_item.setFont(0, bold)
d_item.setForeground(0, domain_color) d_item.setForeground(0, domain_color)
for obj_name, fields in domain_data.items(): for obj_name, fields in domain_data.items():
path = f"{domain}/{obj_name}"
if _is_phys(fields): if _is_phys(fields):
item = QTreeWidgetItem([obj_name, str(fields["value"]), fields["unit"]]) item = QTreeWidgetItem([obj_name, str(fields["value"]), fields["unit"], fields.get("str_state", "")])
item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter) item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter)
if path in changed:
for col in range(4):
item.setBackground(col, _CHANGED_BG)
d_item.addChild(item) d_item.addChild(item)
elif isinstance(fields, dict): elif isinstance(fields, dict):
obj_item = QTreeWidgetItem([obj_name]) obj_item = QTreeWidgetItem([obj_name])
obj_item.setFont(0, bold) obj_item.setFont(0, bold)
self._add_fields(obj_item, fields, f"{domain}/{obj_name}") self._add_fields(obj_item, fields, path, changed)
d_item.addChild(obj_item) d_item.addChild(obj_item)
else: else:
item = QTreeWidgetItem([obj_name, str(fields), ""]) item = QTreeWidgetItem([obj_name, str(fields), "", ""])
item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter) item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter)
if path in changed:
for col in range(4):
item.setBackground(col, _CHANGED_BG)
d_item.addChild(item) d_item.addChild(item)
self._tree.addTopLevelItem(d_item) self._tree.addTopLevelItem(d_item)
self._tree.expandAll() self._tree.expandAll()
self._restore_expanded(expanded) self._apply_collapsed(collapsed)
def _add_fields(self, parent: QTreeWidgetItem, data: dict, path: str): def _add_fields(self, parent: QTreeWidgetItem, data: dict, path: str, changed: set):
for k, v in data.items(): for k, v in data.items():
child_path = f"{path}/{k}" child_path = f"{path}/{k}"
if _is_phys(v): if _is_phys(v):
item = QTreeWidgetItem([k, str(v["value"]), v["unit"]]) item = QTreeWidgetItem([k, str(v["value"]), v["unit"], v.get("str_state", "")])
item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter) item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter)
if child_path in changed:
for col in range(4):
item.setBackground(col, _CHANGED_BG)
parent.addChild(item) parent.addChild(item)
elif isinstance(v, dict): elif isinstance(v, dict):
node = QTreeWidgetItem([k]) node = QTreeWidgetItem([k])
self._add_fields(node, v, child_path) self._add_fields(node, v, child_path, changed)
parent.addChild(node) parent.addChild(node)
else: else:
item = QTreeWidgetItem([k, str(v), ""]) item = QTreeWidgetItem([k, str(v), "", ""])
item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter) item.setTextAlignment(1, Qt.AlignRight | Qt.AlignVCenter)
if child_path in changed:
for col in range(4):
item.setBackground(col, _CHANGED_BG)
parent.addChild(item) parent.addChild(item)
# ── preserve expand state across refreshes ───────────────────────────── # ── preserve collapse state across refreshes ─────────────────────────────
def _expanded_paths(self) -> set[str]: def _collapsed_paths(self) -> set[str]:
paths = set() paths = set()
root = self._tree.invisibleRootItem() self._collect_collapsed(self._tree.invisibleRootItem(), "", paths)
self._collect_expanded(root, "", paths)
return paths return paths
def _collect_expanded(self, item, prefix, paths): def _collect_collapsed(self, item, prefix, paths):
for i in range(item.childCount()): for i in range(item.childCount()):
child = item.child(i) child = item.child(i)
path = f"{prefix}/{child.text(0)}" path = f"{prefix}/{child.text(0)}"
if child.isExpanded(): if child.childCount() > 0 and not child.isExpanded():
paths.add(path) paths.add(path)
self._collect_expanded(child, path, paths) self._collect_collapsed(child, path, paths)
def _restore_expanded(self, paths: set[str]): def _apply_collapsed(self, paths: set[str]):
root = self._tree.invisibleRootItem() self._set_collapsed(self._tree.invisibleRootItem(), "", paths)
self._apply_expanded(root, "", paths)
def _apply_expanded(self, item, prefix, paths): def _set_collapsed(self, item, prefix, paths):
for i in range(item.childCount()): for i in range(item.childCount()):
child = item.child(i) child = item.child(i)
path = f"{prefix}/{child.text(0)}" path = f"{prefix}/{child.text(0)}"
if path in paths: if path in paths:
child.setExpanded(True) child.setExpanded(False)
self._apply_expanded(child, path, paths) self._set_collapsed(child, path, paths)
# ── Plot tab ────────────────────────────────────────────────────────────────── # ── Plot tab ──────────────────────────────────────────────────────────────────
@@ -844,7 +864,7 @@ class MainWindow(QMainWindow):
def _connect(self, host: str, port: int): def _connect(self, host: str, port: int):
self._plot.clear_data() self._plot.clear_data()
err = self._reader.connect_to(host, port, self._connector.diff_mode) err = self._reader.connect_to(host, port)
if err: if err:
self._connector.set_disconnected(err) self._connector.set_disconnected(err)
self._statusbar.showMessage(f"Connection failed: {err}") self._statusbar.showMessage(f"Connection failed: {err}")
+10 -10
View File
@@ -1,20 +1,20 @@
#!/usr/bin/env bash #!/usr/bin/env bash
# CLIENT installer — sets up the we_monitor GUI client. # CLIENT installer — sets up the we_monitor GUI client.
# Does NOT install the server (see server/install.sh). # Does NOT install the server (see install_server.sh).
# #
# Fresh install: # Fresh install:
# chmod +x install.sh # chmod +x install_client.sh
# ./install.sh # ./install_client.sh
# #
# Update (code already pulled via git): # Update (code already pulled via git):
# ./install.sh # ./install_client.sh
set -euo pipefail set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
REPO_DIR="$(dirname "$SCRIPT_DIR")" REPO_DIR="$SCRIPT_DIR"
VENV="$SCRIPT_DIR/.venv" VENV="$REPO_DIR/.venv"
PYTHON="$VENV/bin/python" PYTHON="$VENV/bin/python"
IS_UPDATE=false IS_UPDATE=false
@@ -41,15 +41,15 @@ fi
echo "==> Installing / updating dependencies" echo "==> Installing / updating dependencies"
"$VENV/bin/pip" install --quiet --upgrade pip "$VENV/bin/pip" install --quiet --upgrade pip
"$VENV/bin/pip" install --quiet --upgrade -r "$SCRIPT_DIR/requirements.txt" "$VENV/bin/pip" install --quiet --upgrade -r "$REPO_DIR/client/requirements.txt"
echo " Python : $PYTHON" echo " Python : $PYTHON"
echo "" echo ""
# ── 2. launcher script ──────────────────────────────────────────────────────── # ── 2. launcher script ────────────────────────────────────────────────────────
LAUNCHER="$SCRIPT_DIR/run.sh" LAUNCHER="$REPO_DIR/client/run.sh"
cat > "$LAUNCHER" <<EOF cat > "$LAUNCHER" <<EOF
#!/usr/bin/env bash #!/usr/bin/env bash
exec "$PYTHON" "$SCRIPT_DIR/gui_client.py" "\$@" exec "$PYTHON" "$REPO_DIR/client/gui_client.py" "\$@"
EOF EOF
chmod +x "$LAUNCHER" chmod +x "$LAUNCHER"
echo "==> Launcher written : $LAUNCHER" echo "==> Launcher written : $LAUNCHER"
@@ -60,7 +60,7 @@ CMD="$BIN_DIR/we_monitor_gui"
mkdir -p "$BIN_DIR" mkdir -p "$BIN_DIR"
cat > "$CMD" <<EOF cat > "$CMD" <<EOF
#!/usr/bin/env bash #!/usr/bin/env bash
exec "$PYTHON" "$SCRIPT_DIR/gui_client.py" "\$@" exec "$PYTHON" "$REPO_DIR/client/gui_client.py" "\$@"
EOF EOF
chmod +x "$CMD" chmod +x "$CMD"
echo "==> System command : $CMD" echo "==> System command : $CMD"
+13 -3
View File
@@ -15,7 +15,7 @@
set -euo pipefail set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
REPO_DIR="$(dirname "$SCRIPT_DIR")" REPO_DIR="$SCRIPT_DIR"
SERVICE_NAME="we_monitor" SERVICE_NAME="we_monitor"
CREDS_DIR="$HOME/.config/we_monitor" CREDS_DIR="$HOME/.config/we_monitor"
CREDS_FILE="$CREDS_DIR/credentials.json" CREDS_FILE="$CREDS_DIR/credentials.json"
@@ -23,7 +23,7 @@ SYSTEMD_DIR="$HOME/.config/systemd/user"
SERVICE_FILE="$SYSTEMD_DIR/$SERVICE_NAME.service" SERVICE_FILE="$SYSTEMD_DIR/$SERVICE_NAME.service"
LOG_DIR="$REPO_DIR/logs" LOG_DIR="$REPO_DIR/logs"
VENV="$SCRIPT_DIR/.venv" VENV="$REPO_DIR/.venv"
PYTHON="$VENV/bin/python" PYTHON="$VENV/bin/python"
# detect fresh install vs update # detect fresh install vs update
@@ -55,7 +55,17 @@ fi
echo "==> Installing / updating dependencies" echo "==> Installing / updating dependencies"
"$VENV/bin/pip" install --quiet --upgrade pip "$VENV/bin/pip" install --quiet --upgrade pip
"$VENV/bin/pip" install --quiet --upgrade -r "$SCRIPT_DIR/requirements.txt" "$VENV/bin/pip" install --quiet --upgrade -r "$REPO_DIR/server/requirements.txt"
echo "==> Applying cariad-hybrid-auth-fix.patch to volkswagencarnet"
SITE_PKG="$("$PYTHON" -c 'import sysconfig; print(sysconfig.get_path("purelib"))')"
PATCH_FILE="$REPO_DIR/cariad-hybrid-auth-fix.patch"
if patch -d "$SITE_PKG" -p1 -N --dry-run --silent < "$PATCH_FILE" 2>/dev/null; then
patch -d "$SITE_PKG" -p1 -N --silent < "$PATCH_FILE"
echo " Patch applied"
else
echo " Patch already applied (skipping)"
fi
echo " Python : $PYTHON" echo " Python : $PYTHON"
echo "" echo ""
+32 -14
View File
@@ -6,10 +6,11 @@ from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from . import log_config, network, we_connect from . import log_config, network, we_connect
from .data_model import ALL_DOMAINS, collect_snapshot, apply_procedural
from .storage_helpers import auto_out_path, load_last_24h_records, load_store, save_store
from jaydiff.diff import diff_full as jay_diff_full from jaydiff.diff import diff_full as jay_diff_full
from .data_model import ALL_DOMAINS, apply_procedural, extract_all
from .storage_helpers import auto_out_path, load_last_24h_records, load_store, save_store
from volkswagencarnet.vw_exceptions import AuthenticationError from volkswagencarnet.vw_exceptions import AuthenticationError
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
@@ -46,10 +47,9 @@ def run(args: argparse.Namespace) -> None:
out = auto_out_path(logs_dir, vin) out = auto_out_path(logs_dir, vin)
push_clients, push_lock, push_store_ref = network.start_push_server(args.host, args.port) push_clients, push_lock, push_store_ref = network.start_push_server(args.host, args.port)
push_store_ref["diff_mode"] = args.diff
push_store_ref["preloaded"] = load_last_24h_records(logs_dir, vin) push_store_ref["preloaded"] = load_last_24h_records(logs_dir, vin)
for _r in push_store_ref["preloaded"]: for _r in push_store_ref["preloaded"]:
if "procedural" not in _r: if _r is not None and "procedural" not in _r:
apply_procedural(_r) apply_procedural(_r)
if push_store_ref["preloaded"]: if push_store_ref["preloaded"]:
log.info("Pre-loaded %d record(s) from the last 24 h", len(push_store_ref["preloaded"])) log.info("Pre-loaded %d record(s) from the last 24 h", len(push_store_ref["preloaded"]))
@@ -57,7 +57,7 @@ def run(args: argparse.Namespace) -> None:
current_day = datetime.now(timezone.utc).date() current_day = datetime.now(timezone.utc).date()
store = load_store(out, vin, args.interval) store = load_store(out, vin, args.interval)
for _r in store["records"]: for _r in store["records"]:
if "procedural" not in _r: if _r is not None and "procedural" not in _r:
apply_procedural(_r) apply_procedural(_r)
push_store_ref["current"] = store push_store_ref["current"] = store
last_broadcast: dict = {} last_broadcast: dict = {}
@@ -84,18 +84,36 @@ def run(args: argparse.Namespace) -> None:
log.info("Midnight UTC rotation → %s", out) log.info("Midnight UTC rotation → %s", out)
try: try:
we_connect.update(wc) ok = we_connect.update(wc)
snapshot = collect_snapshot(vehicle, domains) # Re-select vehicle after every update: if doLogin() was triggered
apply_procedural(snapshot) # (rare token-expiry full re-login), the library rebuilds its vehicle
# list with new objects, making our old reference permanently stale.
vehicle, err = we_connect.select_vehicle(wc, args.vin)
if err:
log.warning("Vehicle re-selection after update failed: %s", err)
continue
log.debug(
"update ok=%s vehicle._states keys: %s",
ok,
list(vehicle.attrs.keys()),
)
snapshot = extract_all(vehicle)
if snapshot is None:
log.warning("extract_all returned None — skipping this cycle")
continue
snapshot = apply_procedural(snapshot)
store["records"].append(snapshot) store["records"].append(snapshot)
save_store(out, store, args.max_records) save_store(out, store, args.max_records)
if args.diff and last_broadcast: payload = jay_diff_full(last_broadcast, snapshot, combine_upd_add=True)
payload = jay_diff_full(last_broadcast, snapshot, combine_upd_add=True)
else:
payload = snapshot
network.broadcast(payload, push_clients, push_lock)
last_broadcast = snapshot last_broadcast = snapshot
log.info("Record #%d saved at %s", len(store["records"]), snapshot["ts"]) network.broadcast(payload, push_clients, push_lock)
changed = [k for k in payload if k != "ts"]
log.info(
"Record #%d at %s changed=%s",
len(store["records"]),
snapshot["ts"],
changed or "none",
)
except KeyboardInterrupt: except KeyboardInterrupt:
raise raise
except AuthenticationError: except AuthenticationError:
+74 -6
View File
@@ -1,8 +1,10 @@
import logging import logging
import re
from collections.abc import Callable from collections.abc import Callable
from datetime import datetime, timezone from datetime import datetime, timezone
from volkswagencarnet.vw_const import Paths from volkswagencarnet.vw_const import Paths
from volkswagencarnet.vw_dashboard import Dashboard
from volkswagencarnet.vw_utilities import find_path from volkswagencarnet.vw_utilities import find_path
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
@@ -283,6 +285,73 @@ def extract_windows(vehicle) -> dict | None:
return result or None return result or None
_UNIT_WHITELIST: frozenset[str] = frozenset({"min", "%", "d", "kW", "°C", "km", "V", "W", "A"})
_UNIT_CONVERSIONS: dict[str, str] = {"d": "days", "°C": "degC"}
def _normalise_unit(unit: str) -> str:
if unit not in _UNIT_WHITELIST:
return ""
return _UNIT_CONVERSIONS.get(unit, unit)
def _unit_from_str_state(str_state: str, state) -> str:
"""Parse the unit suffix from a str_state like '42 %' or '21.5 degC'."""
state_repr = str(state)
if str_state.startswith(state_repr):
suffix = str_state[len(state_repr):].strip()
if suffix:
return suffix
# fallback: last whitespace-separated token if it looks like a unit
parts = str_state.rsplit(None, 1)
if len(parts) == 2 and re.fullmatch(r"[^\d\s]\S*", parts[1]):
return parts[1]
return ""
def extract_all(vehicle) -> dict | None:
"""Read every supported instrument via the Dashboard and return a flat snapshot.
Structure: snapshot[component][attr] = {"value": state, "unit": unit, "str_state": str_state}
Unit is taken from instrument.unit when available, otherwise parsed from str_state.
"""
snapshot: dict = {"ts": datetime.now(timezone.utc).isoformat()}
try:
dashboard = Dashboard(vehicle)
except Exception:
log.debug("extract_all: Dashboard setup failed", exc_info=True)
return None
for instrument in dashboard.instruments:
try:
state = instrument.state
if state is None:
continue
try:
str_state = str(instrument.str_state)
except Exception:
str_state = str(state)
if hasattr(instrument, "unit") and instrument.unit:
raw_unit = instrument.unit
else:
raw_unit = _unit_from_str_state(str_state, state)
unit = _normalise_unit(raw_unit)
component = instrument.component
attr = instrument.attr
snapshot.setdefault(component, {})[attr] = {
"value": state,
"unit": unit,
"str_state": str_state,
}
except Exception:
log.debug("extract_all: %s unavailable", instrument.attr, exc_info=True)
return snapshot or None
ALL_DOMAINS: dict[str, Callable] = { ALL_DOMAINS: dict[str, Callable] = {
"charging": extract_charging, "charging": extract_charging,
"climatisation": extract_climatisation, "climatisation": extract_climatisation,
@@ -324,18 +393,17 @@ def apply_procedural(record: dict) -> dict:
""" """
proc: dict = {} proc: dict = {}
range_at_soc = _get(record, "electric_drive", "range", "value") range_at_soc = _get(record, "sensor", "electric_range", "value")
soc = _get(record, "electric_drive", "battery", "soc", "value") soc = _get(record, "sensor", "battery_level", "value")
if range_at_soc is not None and soc: if range_at_soc is not None and soc:
proc["range_at_100"] = _phys(round(100 * float(range_at_soc) / float(soc), 1), "km") proc["range_at_100"] = _phys(round(100 * float(range_at_soc) / float(soc), 1), "km")
# ── add computed fields here ────────────────────────────────────────── # ── add computed fields here ──────────────────────────────────────────
# Use _get(record, "domain", "statusObject", "field", "value") to # Use _get(record, "sensor"|"binary_sensor"|..., "attr_name", "value")
# safely read any nested value. Always use {"value": ..., "unit": ...} # Always use {"value": ..., "unit": ...} format so the GUI picks it up.
# format so the GUI picks up the field automatically.
# #
# Example — usable energy estimated from SOC and battery capacity: # Example — usable energy estimated from SOC and battery capacity:
# soc = _get(record, "electric_drive", "battery", "soc", "value") # soc = _get(record, "sensor", "battery_level", "value")
# if soc is not None: # if soc is not None:
# proc["energy_stored"] = {"value": round(soc / 100 * 77.0, 1), "unit": "kWh"} # proc["energy_stored"] = {"value": round(soc / 100 * 77.0, 1), "unit": "kWh"}
# ───────────────────────────────────────────────────────────────────── # ─────────────────────────────────────────────────────────────────────
-9
View File
@@ -120,12 +120,6 @@ def main() -> None:
help="Directory for rotating log files (overrides credentials file, default: ./logs)", help="Directory for rotating log files (overrides credentials file, default: ./logs)",
) )
p.add_argument("-v", "--verbose", action="store_true", help="Enable debug logging") p.add_argument("-v", "--verbose", action="store_true", help="Enable debug logging")
p.add_argument(
"--diff",
action="store_true",
default=None,
help="Send diffs instead of full snapshots over the push server (overrides credentials file, default: off)",
)
p.add_argument( p.add_argument(
"--dry", "--dry",
action="store_true", action="store_true",
@@ -176,9 +170,6 @@ def main() -> None:
args.port = int(creds_data.get("port", 9999)) args.port = int(creds_data.get("port", 9999))
if args.log_dir is None: if args.log_dir is None:
args.log_dir = creds_data.get("log_dir", "./logs") args.log_dir = creds_data.get("log_dir", "./logs")
if args.diff is None:
args.diff = bool(creds_data.get("diff", False))
if not args.dry and (not args.username or not args.password): if not args.dry and (not args.username or not args.password):
p.error( p.error(
"Username and password are required. " "Username and password are required. "
+16 -14
View File
@@ -24,28 +24,30 @@ def _accept_loop(
preloaded = store_ref.get("preloaded", []) preloaded = store_ref.get("preloaded", [])
if preloaded: if preloaded:
seen_ts = {r.get("ts") for r in preloaded} seen_ts = {r.get("ts") for r in preloaded if r is not None}
extra = [r for r in in_mem if r.get("ts") not in seen_ts] extra = [r for r in in_mem if r is not None and r.get("ts") not in seen_ts]
records = sorted(preloaded + extra, key=lambda r: r.get("ts", "")) records = sorted(preloaded + extra, key=lambda r: r.get("ts", "") if r is not None else "")
else: else:
records = in_mem records = [r for r in in_mem if r is not None]
try: try:
if store_ref.get("diff_mode"): prev: dict = {}
prev: dict = {} for record in records:
for record in records: if record is None:
diff = jay_diff_full(prev, record, combine_upd_add=True) continue
prev = record diff = jay_diff_full(prev, record, combine_upd_add=True)
conn.sendall((json.dumps(diff, default=str) + "\n").encode()) prev = record
else: conn.sendall((json.dumps(diff, default=str) + "\n").encode())
for record in records:
conn.sendall((json.dumps(record, default=str) + "\n").encode())
if records: if records:
log.info("Sent %d record(s) to %s", len(records), addr) log.info("Sent %d record(s) to %s", len(records), addr)
except OSError: except OSError:
log.debug("Failed to send history to %s", addr) log.debug("Failed to send history to %s", addr)
conn.close() conn.close()
continue continue
except Exception:
log.exception("Error sending history to %s — dropping client", addr)
conn.close()
continue
clients.append(conn) clients.append(conn)
except OSError: except OSError:
break break
@@ -73,7 +75,7 @@ def start_push_server(host: str, port: int) -> tuple:
server_sock.listen() server_sock.listen()
clients: list = [] clients: list = []
lock = threading.Lock() lock = threading.Lock()
store_ref: dict = {"current": None, "preloaded": [], "diff_mode": False} store_ref: dict = {"current": None, "preloaded": []}
t = threading.Thread( t = threading.Thread(
target=_accept_loop, target=_accept_loop,
args=(server_sock, clients, lock, store_ref), args=(server_sock, clients, lock, store_ref),
+1 -1
View File
@@ -1,2 +1,2 @@
volkswagencarnet @ git+file:///home/jens/work/repos/volkswagencarnet volkswagencarnet @ git+https://github.com/robinostlund/volkswagencarnet.git@c30dc37
jaydiff @ git+http://192.168.22.90:3001/jayfield/jaypy.git jaydiff @ git+http://192.168.22.90:3001/jayfield/jaypy.git
+2 -2
View File
@@ -53,11 +53,11 @@ def load_last_24h_records(logs_dir: Path, vin: str) -> list:
try: try:
data = json.loads(path.read_text()) data = json.loads(path.read_text())
if isinstance(data, dict) and "records" in data: if isinstance(data, dict) and "records" in data:
kept = [r for r in data["records"] if r.get("ts", "") >= cutoff_str] kept = [r for r in data["records"] if r is not None and r.get("ts", "") >= cutoff_str]
records.extend(kept) records.extend(kept)
log.debug("Loaded %d/%d records from %s", len(kept), len(data["records"]), path.name) log.debug("Loaded %d/%d records from %s", len(kept), len(data["records"]), path.name)
except (json.JSONDecodeError, OSError): except (json.JSONDecodeError, OSError):
log.warning("Could not load %s", path) log.warning("Could not load %s", path)
records.sort(key=lambda r: r.get("ts", "")) records.sort(key=lambda r: r.get("ts", "") if r is not None else "")
return records return records
+8 -2
View File
@@ -60,5 +60,11 @@ def select_vehicle(wc: _Session, vin: str | None):
return vehicles[0], None return vehicles[0], None
def update(wc: _Session) -> None: def update(wc: _Session) -> bool:
wc.run(wc.connection.update()) """Run one update cycle. Returns True if the library reported success."""
ok = wc.run(wc.connection.update())
if not ok:
log.warning("connection.update() returned False — data may not have refreshed")
else:
log.debug("connection.update() succeeded")
return bool(ok)