From ce7dc94241bcf4eaea979725bf7dca83c128067d Mon Sep 17 00:00:00 2001 From: virtuscyber Date: Sat, 11 Apr 2026 11:35:35 -0400 Subject: [PATCH] feat(obd_ai): add deterministic basic health scan workflow --- obd_ai/catalog.py | 21 ++ obd_ai/serializers.py | 4 + obd_ai/session.py | 21 +- obd_ai/tools.py | 374 ++++++++++++++++++++++++++++++- tests/test_obd_ai_catalog.py | 8 + tests/test_obd_ai_serializers.py | 4 + tests/test_obd_ai_session.py | 10 + tests/test_obd_ai_tools.py | 153 ++++++++++++- 8 files changed, 592 insertions(+), 3 deletions(-) diff --git a/obd_ai/catalog.py b/obd_ai/catalog.py index 2a79ed65..f7792a6c 100644 --- a/obd_ai/catalog.py +++ b/obd_ai/catalog.py @@ -150,6 +150,27 @@ def as_mapping(self) -> Mapping[str, ApprovedCommand]: description="Read current ECU/control module voltage.", obd_command_name="CONTROL_MODULE_VOLTAGE", ), + _CommandDefinition( + key="short_fuel_trim_bank_1", + public_name="Short Fuel Trim Bank 1", + category="fuel", + description="Read short-term fuel trim for bank 1.", + obd_command_name="SHORT_FUEL_TRIM_1", + ), + _CommandDefinition( + key="long_fuel_trim_bank_1", + public_name="Long Fuel Trim Bank 1", + category="fuel", + description="Read long-term fuel trim for bank 1.", + obd_command_name="LONG_FUEL_TRIM_1", + ), + _CommandDefinition( + key="mass_air_flow", + public_name="Mass Air Flow", + category="powertrain", + description="Read current mass air flow rate.", + obd_command_name="MAF", + ), _CommandDefinition( key="engine_run_time", public_name="Engine Run Time", diff --git a/obd_ai/serializers.py b/obd_ai/serializers.py index e687cae3..efff11b8 100644 --- a/obd_ai/serializers.py +++ b/obd_ai/serializers.py @@ -146,6 +146,8 @@ def serialize_connection_metadata( status: str, is_connected: Optional[bool] = None, port_name: Optional[str] = None, + protocol_id: Optional[str] = None, + protocol_name: Optional[str] = None, ) -> Dict[str, JsonValue]: """Serialize transport-level connection metadata.""" @@ -154,6 +156,8 @@ def serialize_connection_metadata( "status": status, "is_connected": is_connected, "port_name": port_name, + "protocol_id": protocol_id, + "protocol_name": protocol_name, } diff --git a/obd_ai/session.py b/obd_ai/session.py index 0d4e093f..ab601b1e 100644 --- a/obd_ai/session.py +++ b/obd_ai/session.py @@ -2,7 +2,7 @@ from __future__ import annotations -from typing import Callable, Optional +from typing import Any, Callable, Optional import obd from obd.OBDResponse import OBDResponse @@ -66,12 +66,31 @@ def port_name(self) -> str: return self._connection.port_name() def connection_metadata(self): + protocol_id = self._safe_optional_connection_call("protocol_id") + protocol_name = self._safe_optional_connection_call("protocol_name") + if protocol_id is not None: + protocol_id = str(protocol_id) + if protocol_name is not None: + protocol_name = str(protocol_name) + return serialize_connection_metadata( status=self.status(), is_connected=self.is_connected(), port_name=self.port_name(), + protocol_id=protocol_id, + protocol_name=protocol_name, ) + def _safe_optional_connection_call(self, method_name: str) -> Optional[Any]: + method = getattr(self._connection, method_name, None) + if method is None: + return None + + try: + return method() + except Exception: + return None + def close(self) -> None: self._connection.close() diff --git a/obd_ai/tools.py b/obd_ai/tools.py index 06f57d2e..dab8cf44 100644 --- a/obd_ai/tools.py +++ b/obd_ai/tools.py @@ -2,7 +2,7 @@ from __future__ import annotations -from typing import Any, Dict, Mapping, Optional, Sequence, Tuple +from typing import Any, Dict, List, Mapping, Optional, Sequence, Tuple from .catalog import ApprovedCommandCatalog, DEFAULT_APPROVED_COMMAND_CATALOG from .serializers import serialize_approved_command_catalog, serialize_obd_response @@ -14,6 +14,20 @@ SerialTimeoutException = () # type: ignore[assignment] +_BASIC_HEALTH_SCAN_LIVE_METRICS = ( + "engine_rpm", + "vehicle_speed", + "engine_coolant_temperature", + "control_module_voltage", + "short_fuel_trim_bank_1", + "long_fuel_trim_bank_1", + "mass_air_flow", + "throttle_position", + "intake_air_temperature", + "fuel_level", +) + + class OBDAIReadOnlyToolSurface: """Constrained, read-only tool interface over an ``OBDAISession``.""" @@ -260,6 +274,49 @@ def get_emissions_monitor_status( }, ) + def basic_health_scan(self, input: Optional[Mapping[str, Any]] = None) -> Dict[str, Any]: + """Run a deterministic baseline read-only health scan.""" + + del input + session, error = self._require_connected(tool="basic_health_scan") + if error is not None: + return error + + status = self._scan_read("status_since_dtc_clear") + dtc_snapshot = { + "stored": self._scan_read("stored_trouble_codes"), + "pending": self._scan_read("pending_trouble_codes"), + "freeze_frame": self._scan_read("freeze_frame_trouble_code"), + } + live_metrics = { + command_key: self._scan_read(command_key) + for command_key in _BASIC_HEALTH_SCAN_LIVE_METRICS + } + + findings, anomalies, next_steps, summary = self._analyze_basic_health_scan( + status=status, + dtc_snapshot=dtc_snapshot, + live_metrics=live_metrics, + ) + + return self._success( + "basic_health_scan", + { + "scan_version": "basic_health_scan.v1", + "connection": session.connection_metadata(), + "snapshot": { + "status": status, + "dtc": dtc_snapshot, + "readiness": status, + "live_metrics": live_metrics, + }, + "findings": findings, + "anomalies": anomalies, + "next_steps": next_steps, + "summary": summary, + }, + ) + def close(self) -> None: """Close the active session if one exists.""" @@ -293,6 +350,321 @@ def _multi_read(self, tool: str, reads: Mapping[str, str]) -> Dict[str, Any]: "data": results, } + def _scan_read(self, command_key: str) -> Dict[str, Any]: + response, error = self._query_approved_command( + tool="basic_health_scan", + command_key=command_key, + ) + if error is not None: + return { + "ok": False, + "command_key": command_key, + "error": error["error"], + } + + return { + "ok": True, + "command_key": command_key, + "response": response, + } + + @classmethod + def _analyze_basic_health_scan( + cls, + status: Mapping[str, Any], + dtc_snapshot: Mapping[str, Mapping[str, Any]], + live_metrics: Mapping[str, Mapping[str, Any]], + ) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]], List[Dict[str, Any]], Dict[str, Any]]: + findings: List[Dict[str, Any]] = [] + anomalies: List[Dict[str, Any]] = [] + next_steps: List[Dict[str, Any]] = [] + + unsupported_count = 0 + unavailable_count = 0 + + status_value = cls._extract_value_payload(status) + mil = None + dtc_count = None + incomplete_monitors: List[str] = [] + + if status_value and status_value.get("type") == "status": + mil = bool(status_value.get("mil")) + dtc_count = int(status_value.get("dtc_count", 0)) + incomplete_monitors = [ + test["name"] + for test in status_value.get("tests", []) + if test.get("available") and not test.get("complete") + ] + + if mil: + anomaly = { + "id": "mil_on", + "severity": "warning", + "category": "diagnostics", + "title": "Malfunction Indicator Lamp is ON", + "details": { + "mil": mil, + "dtc_count": dtc_count, + }, + } + findings.append(anomaly) + anomalies.append(anomaly) + + if incomplete_monitors: + findings.append( + { + "id": "incomplete_emissions_monitors", + "severity": "info", + "category": "readiness", + "title": "One or more emissions monitors are incomplete", + "details": { + "count": len(incomplete_monitors), + "monitors": incomplete_monitors, + }, + } + ) + else: + findings.append( + { + "id": "status_unavailable", + "severity": "warning", + "category": "readiness", + "title": "Unable to read monitor status", + "details": {"status_read": dict(status)}, + } + ) + + stored_codes = cls._extract_dtc_codes(dtc_snapshot.get("stored", {})) + pending_codes = cls._extract_dtc_codes(dtc_snapshot.get("pending", {})) + freeze_frame_code = cls._extract_dtc_code(dtc_snapshot.get("freeze_frame", {})) + + if stored_codes: + anomaly = { + "id": "stored_dtcs_present", + "severity": "warning", + "category": "diagnostics", + "title": "Stored DTCs detected", + "details": { + "count": len(stored_codes), + "codes": stored_codes, + }, + } + findings.append(anomaly) + anomalies.append(anomaly) + + if pending_codes: + findings.append( + { + "id": "pending_dtcs_present", + "severity": "info", + "category": "diagnostics", + "title": "Pending DTCs detected", + "details": { + "count": len(pending_codes), + "codes": pending_codes, + }, + } + ) + + if freeze_frame_code is not None: + findings.append( + { + "id": "freeze_frame_dtc", + "severity": "info", + "category": "diagnostics", + "title": "Freeze-frame trigger DTC reported", + "details": freeze_frame_code, + } + ) + + for command_key, metric in live_metrics.items(): + if not metric.get("ok"): + error_code = metric.get("error", {}).get("code") + if error_code == "unsupported_pid": + unsupported_count += 1 + elif error_code in {"null_response", "timeout"}: + unavailable_count += 1 + continue + + magnitude = cls._extract_quantity_magnitude(metric) + if magnitude is None: + continue + + if command_key == "control_module_voltage" and (magnitude < 11.8 or magnitude > 15.0): + anomaly = { + "id": "control_module_voltage_out_of_range", + "severity": "warning", + "category": "electrical", + "title": "Control module voltage is outside expected range", + "details": { + "value": magnitude, + "expected_range": [11.8, 15.0], + }, + } + findings.append(anomaly) + anomalies.append(anomaly) + + if command_key == "engine_coolant_temperature" and magnitude >= 110.0: + anomaly = { + "id": "high_coolant_temperature", + "severity": "warning", + "category": "powertrain", + "title": "Coolant temperature appears high", + "details": { + "value": magnitude, + "threshold": 110.0, + }, + } + findings.append(anomaly) + anomalies.append(anomaly) + + if command_key in {"short_fuel_trim_bank_1", "long_fuel_trim_bank_1"} and abs(magnitude) >= 15.0: + anomaly = { + "id": f"fuel_trim_out_of_range_{command_key}", + "severity": "warning", + "category": "fuel", + "title": "Fuel trim deviates beyond expected range", + "details": { + "command_key": command_key, + "value": magnitude, + "threshold": 15.0, + }, + } + findings.append(anomaly) + anomalies.append(anomaly) + + if unsupported_count: + findings.append( + { + "id": "partially_supported_vehicle", + "severity": "info", + "category": "capability", + "title": "Some approved commands are not supported by this vehicle", + "details": {"unsupported_command_count": unsupported_count}, + } + ) + + if unavailable_count: + findings.append( + { + "id": "temporarily_unavailable_data", + "severity": "info", + "category": "capability", + "title": "Some commands returned no data or timed out", + "details": {"unavailable_command_count": unavailable_count}, + } + ) + + if mil: + next_steps.append( + { + "priority": "high", + "action": "Investigate active MIL and stored DTCs before clearing codes.", + "reason": "MIL indicates an active fault condition.", + } + ) + + if stored_codes: + next_steps.append( + { + "priority": "high", + "action": "Diagnose stored DTC root causes using service information and pinpoint tests.", + "reason": "Stored DTCs represent confirmed faults.", + } + ) + + if pending_codes: + next_steps.append( + { + "priority": "medium", + "action": "Recheck pending DTCs after a complete drive cycle.", + "reason": "Pending DTCs may become confirmed or clear on subsequent cycles.", + } + ) + + if incomplete_monitors: + next_steps.append( + { + "priority": "medium", + "action": "Complete the OEM drive cycle to finish readiness monitors.", + "reason": "Incomplete monitors can hide emissions-related faults.", + } + ) + + if unsupported_count: + next_steps.append( + { + "priority": "low", + "action": "Use OEM-specific diagnostics if additional PIDs are needed.", + "reason": "Vehicle does not support all generic approved commands.", + } + ) + + if not next_steps: + next_steps.append( + { + "priority": "low", + "action": "No immediate action required. Continue normal monitoring.", + "reason": "No fault anomalies were detected in this baseline scan.", + } + ) + + overall_status = "ok" + if anomalies: + overall_status = "attention" + + summary = { + "overall_status": overall_status, + "mil": mil, + "reported_dtc_count": dtc_count, + "stored_dtc_count": len(stored_codes), + "pending_dtc_count": len(pending_codes), + "incomplete_monitor_count": len(incomplete_monitors), + "unsupported_command_count": unsupported_count, + "unavailable_command_count": unavailable_count, + } + + return findings, anomalies, next_steps, summary + + @staticmethod + def _extract_value_payload(read_result: Mapping[str, Any]) -> Optional[Mapping[str, Any]]: + if not read_result.get("ok"): + return None + + response = read_result.get("response", {}) + value = response.get("value") if isinstance(response, Mapping) else None + return value if isinstance(value, Mapping) else None + + @classmethod + def _extract_dtc_codes(cls, read_result: Mapping[str, Any]) -> List[Dict[str, Any]]: + value = cls._extract_value_payload(read_result) + if not value or value.get("type") != "dtc_list": + return [] + + codes = value.get("codes", []) + if not isinstance(codes, list): + return [] + + return [dict(code) for code in codes if isinstance(code, Mapping)] + + @classmethod + def _extract_dtc_code(cls, read_result: Mapping[str, Any]) -> Optional[Dict[str, Any]]: + value = cls._extract_value_payload(read_result) + if not value or value.get("type") != "dtc": + return None + return dict(value) + + @classmethod + def _extract_quantity_magnitude(cls, read_result: Mapping[str, Any]) -> Optional[float]: + value = cls._extract_value_payload(read_result) + if not value or value.get("type") != "quantity": + return None + + magnitude = value.get("magnitude") + if isinstance(magnitude, (int, float)): + return float(magnitude) + return None + def _query_approved_command( self, tool: str, diff --git a/tests/test_obd_ai_catalog.py b/tests/test_obd_ai_catalog.py index 5edddf63..2d8c977d 100644 --- a/tests/test_obd_ai_catalog.py +++ b/tests/test_obd_ai_catalog.py @@ -33,3 +33,11 @@ def test_catalog_mapping_is_read_only(): with pytest.raises(TypeError): mapping["new"] = "value" + + +def test_catalog_includes_health_scan_live_metric_commands(): + catalog = DEFAULT_APPROVED_COMMAND_CATALOG + + assert "mass_air_flow" in catalog + assert "short_fuel_trim_bank_1" in catalog + assert "long_fuel_trim_bank_1" in catalog diff --git a/tests/test_obd_ai_serializers.py b/tests/test_obd_ai_serializers.py index a79543fe..664395df 100644 --- a/tests/test_obd_ai_serializers.py +++ b/tests/test_obd_ai_serializers.py @@ -110,6 +110,8 @@ def test_serialize_connection_metadata(): status="Car Connected", is_connected=True, port_name="/dev/ttyUSB0", + protocol_id="6", + protocol_name="ISO 15765-4 (CAN 11/500)", ) assert payload == { @@ -117,4 +119,6 @@ def test_serialize_connection_metadata(): "status": "Car Connected", "is_connected": True, "port_name": "/dev/ttyUSB0", + "protocol_id": "6", + "protocol_name": "ISO 15765-4 (CAN 11/500)", } diff --git a/tests/test_obd_ai_session.py b/tests/test_obd_ai_session.py index a741659f..37c2cdab 100644 --- a/tests/test_obd_ai_session.py +++ b/tests/test_obd_ai_session.py @@ -32,6 +32,14 @@ def port_name(): def supports(command): return command is obd.commands.RPM + @staticmethod + def protocol_id(): + return "6" + + @staticmethod + def protocol_name(): + return "ISO 15765-4 (CAN 11/500)" + def close(self): self.closed = True @@ -116,6 +124,8 @@ def test_session_exposes_connection_metadata(): "status": OBDStatus.CAR_CONNECTED, "is_connected": True, "port_name": "FAKEPORT", + "protocol_id": "6", + "protocol_name": "ISO 15765-4 (CAN 11/500)", } diff --git a/tests/test_obd_ai_tools.py b/tests/test_obd_ai_tools.py index e33a83cb..7dc92bd9 100644 --- a/tests/test_obd_ai_tools.py +++ b/tests/test_obd_ai_tools.py @@ -1,5 +1,5 @@ import obd -from obd.OBDResponse import OBDResponse, Status +from obd.OBDResponse import OBDResponse, Status, StatusTest from obd.utils import OBDStatus from obd_ai.session import OBDAISessionManager @@ -48,6 +48,14 @@ def port_name(): def supports(self, command): return self.support_overrides.get(command.name, True) + @staticmethod + def protocol_id(): + return "6" + + @staticmethod + def protocol_name(): + return "ISO 15765-4 (CAN 11/500)" + def close(self): self.closed = True @@ -70,6 +78,18 @@ def _null_response(command): return OBDResponse(command=command, messages=[]) +def _quantity_response(command, magnitude, unit): + response = OBDResponse(command=command, messages=[object()]) + response.value = obd.Unit.Quantity(magnitude, unit) + return response + + +def _dtc_list_response(command, codes): + response = OBDResponse(command=command, messages=[object()]) + response.value = list(codes) + return response + + def test_connect_vehicle_returns_success_payload_for_connected_adapter(): surface = _surface_for(FakeConnection(connected=True)) @@ -213,3 +233,134 @@ def test_get_freeze_frame_and_emissions_monitor_status_use_approved_reads(): assert status_payload["ok"] is True assert status_payload["data"]["response"]["value"]["type"] == "status" + + +def test_basic_health_scan_returns_structured_snapshot_for_nominal_vehicle(): + status = Status() + status.MIL = False + status.DTC_count = 0 + status.ignition_type = "spark" + status.MISFIRE_MONITORING = StatusTest("MISFIRE_MONITORING", True, True) + + response_overrides = { + "STATUS": OBDResponse(command=obd.commands.STATUS, messages=[object()]), + "GET_DTC": _dtc_list_response(obd.commands.GET_DTC, []), + "GET_CURRENT_DTC": _dtc_list_response(obd.commands.GET_CURRENT_DTC, []), + "FREEZE_DTC": _null_response(obd.commands.FREEZE_DTC), + "RPM": _quantity_response(obd.commands.RPM, 780, obd.Unit.rpm), + "SPEED": _quantity_response(obd.commands.SPEED, 0, obd.Unit.kph), + "COOLANT_TEMP": _quantity_response(obd.commands.COOLANT_TEMP, 89, obd.Unit.celsius), + "CONTROL_MODULE_VOLTAGE": _quantity_response( + obd.commands.CONTROL_MODULE_VOLTAGE, + 13.9, + obd.Unit.volt, + ), + "SHORT_FUEL_TRIM_1": _quantity_response( + obd.commands.SHORT_FUEL_TRIM_1, + 2.5, + obd.Unit.percent, + ), + "LONG_FUEL_TRIM_1": _quantity_response( + obd.commands.LONG_FUEL_TRIM_1, + -1.5, + obd.Unit.percent, + ), + "MAF": _quantity_response(obd.commands.MAF, 4.8, obd.Unit.gram / obd.Unit.second), + "THROTTLE_POS": _quantity_response(obd.commands.THROTTLE_POS, 16, obd.Unit.percent), + "INTAKE_TEMP": _quantity_response(obd.commands.INTAKE_TEMP, 28, obd.Unit.celsius), + "FUEL_LEVEL": _quantity_response(obd.commands.FUEL_LEVEL, 62, obd.Unit.percent), + } + response_overrides["STATUS"].value = status + + connection = FakeConnection(response_overrides=response_overrides) + surface = _surface_for(connection) + surface.connect_vehicle() + + payload = surface.basic_health_scan() + + assert payload["ok"] is True + assert payload["tool"] == "basic_health_scan" + assert payload["data"]["scan_version"] == "basic_health_scan.v1" + assert payload["data"]["summary"]["overall_status"] == "ok" + assert payload["data"]["summary"]["stored_dtc_count"] == 0 + assert payload["data"]["summary"]["pending_dtc_count"] == 0 + assert payload["data"]["summary"]["unsupported_command_count"] == 0 + assert payload["data"]["snapshot"]["live_metrics"]["engine_rpm"]["ok"] is True + assert payload["data"]["snapshot"]["live_metrics"]["mass_air_flow"]["ok"] is True + assert payload["data"]["anomalies"] == [] + + +def test_basic_health_scan_reports_faults_and_partial_support(): + status = Status() + status.MIL = True + status.DTC_count = 2 + status.ignition_type = "spark" + status.MISFIRE_MONITORING = StatusTest("MISFIRE_MONITORING", True, False) + + stored_codes = [ + ("P0300", "Random/Multiple Cylinder Misfire Detected"), + ("P0171", "System Too Lean (Bank 1)"), + ] + + freeze = OBDResponse(command=obd.commands.FREEZE_DTC, messages=[object()]) + freeze.value = ("P0300", "Random/Multiple Cylinder Misfire Detected") + + response_overrides = { + "STATUS": OBDResponse(command=obd.commands.STATUS, messages=[object()]), + "GET_DTC": _dtc_list_response(obd.commands.GET_DTC, stored_codes), + "GET_CURRENT_DTC": _dtc_list_response( + obd.commands.GET_CURRENT_DTC, + [("P0171", "System Too Lean (Bank 1)")], + ), + "FREEZE_DTC": freeze, + "RPM": _quantity_response(obd.commands.RPM, 820, obd.Unit.rpm), + "SPEED": _quantity_response(obd.commands.SPEED, 0, obd.Unit.kph), + "COOLANT_TEMP": _quantity_response(obd.commands.COOLANT_TEMP, 113, obd.Unit.celsius), + "CONTROL_MODULE_VOLTAGE": _quantity_response( + obd.commands.CONTROL_MODULE_VOLTAGE, + 11.2, + obd.Unit.volt, + ), + "SHORT_FUEL_TRIM_1": _quantity_response( + obd.commands.SHORT_FUEL_TRIM_1, + 19, + obd.Unit.percent, + ), + "LONG_FUEL_TRIM_1": _quantity_response( + obd.commands.LONG_FUEL_TRIM_1, + 22, + obd.Unit.percent, + ), + "THROTTLE_POS": _quantity_response(obd.commands.THROTTLE_POS, 20, obd.Unit.percent), + "INTAKE_TEMP": _quantity_response(obd.commands.INTAKE_TEMP, 36, obd.Unit.celsius), + } + response_overrides["STATUS"].value = status + + connection = FakeConnection( + response_overrides=response_overrides, + support_overrides={ + "FUEL_LEVEL": False, + "MAF": False, + }, + ) + surface = _surface_for(connection) + surface.connect_vehicle() + + payload = surface.basic_health_scan() + + assert payload["ok"] is True + assert payload["data"]["summary"]["overall_status"] == "attention" + assert payload["data"]["summary"]["mil"] is True + assert payload["data"]["summary"]["stored_dtc_count"] == 2 + assert payload["data"]["summary"]["pending_dtc_count"] == 1 + assert payload["data"]["summary"]["unsupported_command_count"] == 2 + + anomaly_ids = {item["id"] for item in payload["data"]["anomalies"]} + assert "mil_on" in anomaly_ids + assert "stored_dtcs_present" in anomaly_ids + assert "high_coolant_temperature" in anomaly_ids + assert "control_module_voltage_out_of_range" in anomaly_ids + + next_step_actions = {item["action"] for item in payload["data"]["next_steps"]} + assert "Investigate active MIL and stored DTCs before clearing codes." in next_step_actions + assert "Diagnose stored DTC root causes using service information and pinpoint tests." in next_step_actions