Source code for ur_client

from __future__ import annotations

import json
import re
import socket
import sqlite3
import struct
import threading
import time
import traceback
import zipfile
from collections import deque
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterator


JOINT_LABELS = ["base", "shoulder", "elbow", "wrist_1", "wrist_2", "wrist_3"]
INTERFACE_PORTS: dict[str, int] = {
    "primary": 30001,
    "secondary": 30002,
    "primary_ro": 30011,
    "secondary_ro": 30012,
    "custom": 0,
}
INTERFACE_LABELS: dict[str, str] = {
    "primary": "Primary (30001) - Robot State + Robot Message + Program State",
    "secondary": "Secondary (30002) - Robot State + Version",
    "primary_ro": "Primary Read-Only (30011) - Robot State + Robot Message",
    "secondary_ro": "Secondary Read-Only (30012) - Robot State + Version",
    "custom": "Custom Port - Manual Entry",
}
INTERFACE_NOTES: dict[str, str] = {
    "primary": "Primary (30001) sends robot state plus additional Robot Message packets. KeyMessage is expected here.",
    "secondary": "Secondary (30002) is mainly robot state plus Version. Do not rely on extra Robot Message packets here.",
    "primary_ro": "Primary Read-Only (30011) also sends robot state and Robot Message packets. KeyMessage may appear here.",
    "secondary_ro": "Secondary Read-Only (30012) is mainly robot state plus Version. Do not rely on extra Robot Message packets here.",
    "custom": "Custom Port follows the behavior of the port you specify manually.",
}
MESSAGE_TYPES = {16: "robot_state", 20: "robot_message", 25: "program_state_message"}
ROBOT_MESSAGE_TYPES = {
    0: "text",
    2: "popup",
    3: "version",
    5: "safety_mode",
    6: "error_code",
    7: "key",
    9: "request_value",
    10: "runtime_exception",
    14: "program_threads",
}
PROGRAM_STATE_MESSAGE_TYPES = {
    0: "global_variables_setup",
    1: "global_variables_update",
    2: "variable_update",
}
PROGRAM_VALUE_TYPES = {
    0: "none",
    3: "const_string",
    4: "var_string",
    5: "list",
    10: "pose",
    12: "bool",
    13: "num",
    14: "int",
    15: "float",
}
ROBOT_MODES = {
    -1: "NO_CONTROLLER",
    0: "DISCONNECTED",
    1: "CONFIRM_SAFETY",
    2: "BOOTING",
    3: "POWER_OFF",
    4: "POWER_ON",
    5: "IDLE",
    6: "BACKDRIVE",
    7: "RUNNING",
    8: "UPDATING_FIRMWARE",
}
CONTROL_MODES = {0: "POSITION", 1: "TEACH", 2: "FORCE", 3: "TORQUE"}
JOINT_MODES = {
    235: "RESET",
    236: "SHUTTING_DOWN",
    237: "PART_D_CALIBRATION",
    238: "BACKDRIVE",
    239: "POWER_OFF",
    240: "READY_FOR_POWER_OFF",
    245: "NOT_RESPONDING",
    246: "MOTOR_INITIALISATION",
    247: "BOOTING",
    249: "BOOTLOADER",
    250: "CALIBRATION",
    252: "FAULT",
    253: "RUNNING",
    255: "IDLE",
}
TOOL_MODES = {
    235: "RESET",
    236: "SHUTTING_DOWN",
    239: "POWER_OFF",
    245: "NOT_RESPONDING",
    247: "BOOTING",
    249: "BOOTLOADER",
    252: "FAULT",
    253: "RUNNING",
    255: "IDLE",
}
SAFETY_MODES = {
    1: "NORMAL",
    2: "REDUCED",
    3: "PROTECTIVE_STOP",
    4: "RECOVERY",
    5: "SAFEGUARD_STOP",
    6: "SYSTEM_EMERGENCY_STOP",
    7: "ROBOT_EMERGENCY_STOP",
    8: "VIOLATION",
    9: "FAULT",
    10: "VALIDATE_JOINT_ID",
    11: "UNDEFINED",
}
REPORT_LEVELS = {
    0: "DEBUG",
    1: "INFO",
    2: "WARNING",
    3: "VIOLATION",
    4: "FAULT",
    128: "DEVL_DEBUG",
    129: "DEVL_INFO",
    130: "DEVL_WARNING",
    131: "DEVL_VIOLATION",
    132: "DEVL_FAULT",
}
ROBOT_MESSAGE_DATA_TYPES = {
    0: "uint32",
    1: "uint32",
    2: "int32",
    3: "float",
    4: "hex_uint32",
    5: "string",
}
UR_ERROR_CODE_PATTERN = re.compile(r"\b([A-Z]\d{1,4}(?:A(?:-?\d{1,4})?)?)\b")
REQUESTED_TYPES = {
    0: "BOOLEAN",
    1: "INTEGER",
    2: "FLOAT",
    3: "STRING",
    4: "POSE",
    5: "JOINTVECTOR",
    12: "LIST",
    15: "UNKNOWN",
}
MESSAGE_SOURCES = {
    -5: "GUI",
    -4: "SIMULATED_ROBOT",
    -3: "RTMACHINE",
    -2: "ROBOTINTERFACE",
    -1: "UNDEFINED",
    0: "JOINT_0 (BASE)",
    1: "JOINT_1 (SHOULDER)",
    2: "JOINT_2 (ELBOW)",
    3: "JOINT_3 (WRIST_1)",
    4: "JOINT_4 (WRIST_2)",
    5: "JOINT_5 (WRIST_3)",
    6: "TOOL",
    7: "CONTROLLER",
    8: "RTDE",
    20: "SAFETY_PROCESSOR_UA",
    30: "SAFETY_PROCESSOR_UB",
    40: "SCB_FPGA",
    65: "TEACH_PENDANT_1",
    66: "TEACH_PENDANT_2",
    67: "EUROMAP_1",
    68: "EUROMAP_2",
    100: "JOINT_0_FPGA",
    101: "JOINT_1_FPGA",
    102: "JOINT_2_FPGA",
    103: "JOINT_3_FPGA",
    104: "JOINT_4_FPGA",
    105: "JOINT_5_FPGA",
    106: "TOOL_FPGA",
    107: "EUROMAP_FPGA",
    108: "TEACH_PENDANT_A",
    110: "JOINT_0_A",
    111: "JOINT_1_A",
    112: "JOINT_2_A",
    113: "JOINT_3_A",
    114: "JOINT_4_A",
    115: "JOINT_5_A",
    116: "TOOL_A",
    117: "EUROMAP_A",
    118: "TEACH_PENDANT_B",
    120: "JOINT_0_B",
    121: "JOINT_1_B",
    122: "JOINT_2_B",
    123: "JOINT_3_B",
    124: "JOINT_4_B",
    125: "JOINT_5_B",
    126: "TOOL_B",
    127: "EUROMAP_B",
}
ROBOT_STATE_SUBPACKAGES = {
    0: "robot_mode_data",
    1: "joint_data",
    2: "tool_data",
    3: "masterboard_data",
    4: "cartesian_info",
    5: "kinematics_info",
    6: "configuration_data",
    7: "force_mode_data",
    8: "additional_info",
    9: "calibration_data",
    10: "safety_data",
    11: "tool_communication_info",
    12: "tool_mode_info",
    13: "singularity_info",
}


class ParseError(RuntimeError):
    pass


class PacketLengthError(RuntimeError):
    pass


[docs] @dataclass class AppSettings: host: str = "192.168.163.128" interface_name: str = "primary_ro" port: int | None = None connect_timeout: float = 5.0 read_timeout: float = 5.0 reconnect_delay: float = 2.0 max_packet_size: int = 16 * 1024 * 1024 state_write_interval: float = 0.5 save_sqlite: bool = True sqlite_path: str = "runtime/ur_monitor.sqlite3" save_ndjson: bool = True ndjson_path: str = "runtime/events.ndjson" save_snapshot_json: bool = True snapshot_json_path: str = "runtime/latest_snapshot.json" diagnostics_enabled: bool = True diagnostics_path: str = "runtime/diagnostics.ndjson" diagnostics_export_dir: str = "runtime/debug_exports" error_db_enabled: bool = True error_db_path: str = "ur_error_codes.sqlite3" recent_packet_limit: int = 400 auto_start: bool = True window_refresh_ms: int = 500 message_history_limit: int = 200 favorite_keys: list[str] = field( default_factory=lambda: [ "robot.mode", "robot.control_mode", "robot.safety_mode", "robot.is_program_running", "robot.is_protective_stopped", "tcp.pose.x", "tcp.pose.y", "tcp.pose.z", "joints.base.q_actual", "joints.shoulder.q_actual", "joints.elbow.q_actual", "joints.wrist_1.q_actual", "joints.wrist_2.q_actual", "joints.wrist_3.q_actual", "messages.latest_key.robot_message_code", "messages.latest_key.robot_message_title", "messages.latest_key.key_text_message", "messages.latest_with_code.robot_message_code", "messages.latest_with_code.robot_message_argument", "messages.latest_with_code.category", "messages.latest_with_code.text", "messages.latest_nonzero_code.robot_message_code", "errors.latest.code", "errors.latest.group_title", "errors.latest.description", "errors.latest.explanation", "errors.latest.suggestion", "errors.latest.source_kind", ] )
@dataclass class ParsedPacket: received_at: str message_type: int message_type_name: str message_subtype: int | None message_subtype_name: str | None controller_timestamp_us: int | None source: int | None payload: dict[str, Any] raw: bytes def event_dict(self) -> dict[str, Any]: return { "received_at": self.received_at, "message_type": self.message_type, "message_type_name": self.message_type_name, "message_subtype": self.message_subtype, "message_subtype_name": self.message_subtype_name, "controller_timestamp_us": self.controller_timestamp_us, "source": self.source, "payload": json_safe(self.payload), "raw_hex": self.raw.hex(), } @dataclass class MessageRecord: received_at: str category: str title: str text: str payload: dict[str, Any] message_type_name: str | None = None message_subtype_name: str | None = None controller_timestamp_us: int | None = None source: int | None = None source_name: str | None = None robot_message_code: int | None = None robot_message_argument: int | None = None formatted_error_code: str | None = None detected_error_code: str | None = None request_id: int | None = None requested_type: Any = None def to_dict(self) -> dict[str, Any]: return { "received_at": self.received_at, "category": self.category, "title": self.title, "text": self.text, "message_type_name": self.message_type_name, "message_subtype_name": self.message_subtype_name, "controller_timestamp_us": self.controller_timestamp_us, "source": self.source, "source_name": self.source_name, "robot_message_code": self.robot_message_code, "robot_message_argument": self.robot_message_argument, "formatted_error_code": self.formatted_error_code, "detected_error_code": self.detected_error_code, "request_id": self.request_id, "requested_type": json_safe(self.requested_type), "payload": json_safe(self.payload), } @dataclass class KeyMessage: received_at: str timestamp: int | None source: int | None source_name: str | None robot_message_code: int | None robot_message_argument: int | None robot_message_title_size: int | None robot_message_title: str key_text_message: str def to_dict(self) -> dict[str, Any]: return { "received_at": self.received_at, "timestamp": self.timestamp, "source": self.source, "source_name": self.source_name, "robot_message_code": self.robot_message_code, "robot_message_argument": self.robot_message_argument, "robot_message_title_size": self.robot_message_title_size, "robot_message_title": self.robot_message_title, "key_text_message": self.key_text_message, } @dataclass class ByteReader: buffer: memoryview offset: int = 0 def __init__(self, data: bytes | bytearray | memoryview) -> None: self.buffer = memoryview(data) self.offset = 0 def remaining(self) -> int: return len(self.buffer) - self.offset def empty(self) -> bool: return self.remaining() == 0 def _require(self, size: int, label: str) -> None: if size < 0: raise ParseError(f"invalid negative size for {label}: {size}") if self.remaining() < size: raise ParseError(f"not enough bytes for {label}: need {size}, have {self.remaining()}") def _read(self, fmt: str, label: str) -> Any: size = struct.calcsize(fmt) self._require(size, label) value = struct.unpack_from(fmt, self.buffer, self.offset)[0] self.offset += size return value def read_uint8(self) -> int: return int(self._read(">B", "uint8")) def read_int8(self) -> int: return int(self._read(">b", "int8")) def read_bool(self) -> bool: return bool(self.read_uint8()) def read_uint16(self) -> int: return int(self._read(">H", "uint16")) def read_int16(self) -> int: return int(self._read(">h", "int16")) def read_uint32(self) -> int: return int(self._read(">I", "uint32")) def read_int32(self) -> int: return int(self._read(">i", "int32")) def read_uint64(self) -> int: return int(self._read(">Q", "uint64")) def read_float(self) -> float: return float(self._read(">f", "float")) def read_double(self) -> float: return float(self._read(">d", "double")) def read_bytes(self, size: int) -> bytes: self._require(size, "bytes") start = self.offset end = self.offset + size self.offset = end return self.buffer[start:end].tobytes() def read_reader(self, size: int) -> "ByteReader": return ByteReader(self.read_bytes(size)) def read_string(self, size: int | None = None, encoding: str = "utf-8") -> str: if size is None: size = self.remaining() return self.read_bytes(size).decode(encoding, errors="replace") def remaining_bytes(self) -> bytes: return self.buffer[self.offset :].tobytes() def skip(self, size: int) -> None: self._require(size, "skip") self.offset += size @dataclass class ParseContext: global_variable_names: list[str | None] = field(default_factory=list) global_variable_values: list[dict[str, Any] | None] = field(default_factory=list) def _ensure_size(self, items: list[Any], size: int) -> None: while len(items) < size: items.append(None) def set_variable_names(self, start_index: int, names: list[str]) -> None: self._ensure_size(self.global_variable_names, start_index + len(names)) for offset, name in enumerate(names): self.global_variable_names[start_index + offset] = name def set_variable_values(self, start_index: int, values: list[dict[str, Any]]) -> None: self._ensure_size(self.global_variable_values, start_index + len(values)) for offset, value in enumerate(values): self.global_variable_values[start_index + offset] = value def named_variable_map(self) -> dict[str, Any]: result: dict[str, Any] = {} size = max(len(self.global_variable_names), len(self.global_variable_values)) for index in range(size): name = self.global_variable_names[index] if index < len(self.global_variable_names) else None value = self.global_variable_values[index] if index < len(self.global_variable_values) else None if not name: continue if isinstance(value, dict) and "value" in value: result[name] = value["value"] else: result[name] = value return result def utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat(timespec="milliseconds") def interface_options() -> list[tuple[str, str]]: return list(INTERFACE_LABELS.items()) def interface_label(name: str) -> str: return INTERFACE_LABELS.get(name, name) def interface_note(name: str) -> str: return INTERFACE_NOTES.get(name, "") def supports_key_messages(interface_name: str) -> bool: return interface_name in {"primary", "primary_ro"} def enum_name(mapping: dict[int, str], value: int, fallback: str) -> str: return mapping.get(value, f"{fallback}({value})") def resolve_port(interface_name: str, port: int | None) -> int: if port is not None: return int(port) if interface_name == "custom": raise ValueError("The custom interface requires a manually specified port.") if interface_name not in INTERFACE_PORTS: raise ValueError(f"unknown interface_name: {interface_name!r}") return INTERFACE_PORTS[interface_name] def ensure_parent(path: str | Path) -> Path: resolved = Path(path).expanduser().resolve() resolved.parent.mkdir(parents=True, exist_ok=True) return resolved def json_safe(value: Any) -> Any: if value is None or isinstance(value, (bool, int, float, str)): return value if isinstance(value, bytes): return value.hex() if isinstance(value, Path): return str(value) if isinstance(value, dict): return {str(k): json_safe(v) for k, v in value.items()} if isinstance(value, (list, tuple, set, deque)): return [json_safe(v) for v in value] if hasattr(value, "__dict__"): return json_safe(vars(value)) return str(value) def deep_copy_jsonable(value: Any) -> Any: return json.loads(json.dumps(json_safe(value), ensure_ascii=False)) def format_value(value: Any) -> str: safe = json_safe(value) if isinstance(safe, float): return f"{safe:.6f}" if isinstance(safe, (dict, list)): return json.dumps(safe, ensure_ascii=False) return "" if safe is None else str(safe) def pretty_json(value: Any) -> str: return json.dumps(json_safe(value), ensure_ascii=False, indent=2) def format_robot_error_code(robot_message_code: Any, robot_message_argument: Any) -> str | None: try: code = int(robot_message_code) arg = int(robot_message_argument) except (TypeError, ValueError): return None if code <= 0 or arg < 0: return None return f"C{code}A{arg}" def detect_ur_error_code(*parts: Any) -> str | None: for part in parts: if part is None: continue match = UR_ERROR_CODE_PATTERN.search(str(part)) if match: return match.group(1) return None def normalize_ur_error_code(value: Any) -> str | None: if value is None: return None match = UR_ERROR_CODE_PATTERN.search(str(value).strip().upper()) if not match: return None return match.group(1) def candidate_error_lookup_codes(normalized_code: str | None) -> list[str]: if normalized_code is None: return [] candidates: list[str] = [] def add(value: str | None) -> None: if value and value not in candidates: candidates.append(value) add(normalized_code) match = re.match(r"^(C\d+)A(-?\d+)$", normalized_code) if match: add(f"{match.group(1)}A") add(match.group(1)) return candidates match = re.match(r"^(C\d+)A$", normalized_code) if match: add(match.group(1)) return candidates def ascii_preview(data: bytes, limit: int = 240) -> str: chunk = data[: max(0, int(limit))] text = "".join(chr(byte) if 32 <= byte <= 126 else "." for byte in chunk) if len(data) > len(chunk): text += "..." return text def extract_printable_strings(data: bytes, min_length: int = 4, max_results: int = 20) -> list[str]: results: list[str] = [] current = bytearray() for byte in data: if 32 <= byte <= 126: current.append(byte) else: if len(current) >= min_length: results.append(current.decode("ascii", errors="replace")) if len(results) >= max_results: return results current.clear() if len(current) >= min_length and len(results) < max_results: results.append(current.decode("ascii", errors="replace")) return results def raw_packet_summary(raw: bytes, received_at: str) -> dict[str, Any]: declared_size: int | None = None message_type: int | None = None message_type_name: str | None = None if len(raw) >= 4: declared_size = struct.unpack_from(">I", raw, 0)[0] if len(raw) >= 5: message_type = raw[4] message_type_name = enum_name(MESSAGE_TYPES, message_type, "unknown_message_type") strings_found = extract_printable_strings(raw) return { "received_at": received_at, "length": len(raw), "declared_size": declared_size, "message_type": message_type, "message_type_name": message_type_name, "ascii_preview": ascii_preview(raw), "strings_found": strings_found, "suspected_error_code": detect_ur_error_code(*strings_found), "raw_hex": raw.hex(), } class ErrorCodeDirectory: def __init__(self, sqlite_path: str | Path | None) -> None: self.path = self._resolve_path(sqlite_path) self.available = False self.last_error = "" self._lock = threading.RLock() self._cache: dict[str, dict[str, Any]] = {} self._con: sqlite3.Connection | None = None if self.path is None: return if not self.path.exists(): self.last_error = f"error db not found: {self.path}" return try: self._con = sqlite3.connect(self.path, check_same_thread=False) self._con.row_factory = sqlite3.Row self.available = True except Exception as exc: self.last_error = str(exc) self.available = False @staticmethod def _resolve_path(sqlite_path: str | Path | None) -> Path | None: if sqlite_path in (None, ""): return None candidate = Path(str(sqlite_path)).expanduser() if candidate.is_absolute(): return candidate.resolve() direct = candidate.resolve() if direct.exists(): return direct module_relative = (Path(__file__).resolve().parent / candidate).resolve() return module_relative def close(self) -> None: with self._lock: if self._con is not None: try: self._con.close() finally: self._con = None def _base_result(self, input_code: Any, normalized_code: str | None, lookup_code: str | None, found: bool, row: dict[str, Any] | None) -> dict[str, Any]: result: dict[str, Any] = { "input_code": "" if input_code is None else str(input_code), "normalized_code": normalized_code, "lookup_code": lookup_code or normalized_code, "error_code": lookup_code or normalized_code, "found": bool(found), "db_available": bool(self.available), "db_path": "" if self.path is None else str(self.path), "db_error": self.last_error or "", } if row: result.update(row) result["error_code"] = row.get("error_code") or row.get("code") or result["error_code"] return result def lookup(self, code: str | None) -> dict[str, Any] | None: normalized = normalize_ur_error_code(code) if normalized is None: return None with self._lock: if normalized in self._cache: return deep_copy_jsonable(self._cache[normalized]) if not self.available or self._con is None: result = self._base_result(code, normalized, normalized, False, None) self._cache[normalized] = deep_copy_jsonable(result) return deep_copy_jsonable(result) try: row = None lookup_code = normalized for candidate in candidate_error_lookup_codes(normalized): row = self._con.execute( "SELECT * FROM error_codes_resolved WHERE code = ?", (candidate,), ).fetchone() lookup_code = candidate if row is not None: break row_dict = dict(row) if row is not None else None result = self._base_result(code, normalized, lookup_code, row is not None, row_dict) except Exception as exc: self.last_error = str(exc) result = self._base_result(code, normalized, normalized, False, None) self._cache[normalized] = deep_copy_jsonable(result) return deep_copy_jsonable(result) def search(self, query: str, limit: int = 20) -> list[dict[str, Any]]: query = str(query or "").strip() if not query: return [] with self._lock: if not self.available or self._con is None: return [] try: rows = self._con.execute( """ SELECT ec.* FROM error_code_search fts JOIN error_codes_resolved ec ON ec.code = fts.code WHERE error_code_search MATCH ? LIMIT ? """, (query, int(limit)), ).fetchall() return [dict(row) for row in rows] except Exception as exc: self.last_error = str(exc) return [] class Diagnostics: def __init__(self, settings: AppSettings) -> None: self.settings = settings self.enabled = bool(settings.diagnostics_enabled) self.path = ensure_parent(settings.diagnostics_path) if self.enabled else None self.export_dir = Path(settings.diagnostics_export_dir).expanduser().resolve() self.export_dir.mkdir(parents=True, exist_ok=True) self._lock = threading.RLock() self._fh = self.path.open("a", encoding="utf-8", buffering=1) if self.path is not None else None self._recent_packets: deque[dict[str, Any]] = deque(maxlen=max(50, int(settings.recent_packet_limit))) self._last_export_path = "" def close(self) -> None: with self._lock: if self._fh is not None: self._fh.close() self._fh = None def reset_files(self) -> None: with self._lock: if self._fh is not None: self._fh.close() self._fh = None if self.path is not None and self.path.exists(): self.path.unlink() self.export_dir.mkdir(parents=True, exist_ok=True) self._recent_packets.clear() if self.path is not None: self._fh = self.path.open("a", encoding="utf-8", buffering=1) self._last_export_path = "" def write(self, kind: str, data: dict[str, Any]) -> None: entry = {"logged_at": utc_now_iso(), "kind": kind, **json_safe(data)} with self._lock: if self._fh is not None: self._fh.write(json.dumps(entry, ensure_ascii=False) + "\n") self._fh.flush() def record_raw_packet(self, raw: bytes, received_at: str) -> dict[str, Any]: summary = raw_packet_summary(raw, received_at) with self._lock: self._recent_packets.append(deep_copy_jsonable(summary)) should_write = ( summary.get("message_type") == 20 or bool(summary.get("suspected_error_code")) or str(summary.get("message_type_name") or "").startswith("unknown_") ) if should_write: self.write("raw_packet", summary) return summary def record_parsed_packet(self, packet: ParsedPacket) -> None: if packet.message_type_name == "robot_message" or str(packet.message_type_name).startswith("unknown_") or str(packet.message_subtype_name or "").startswith("unknown_"): self.write( "parsed_packet", { "received_at": packet.received_at, "message_type": packet.message_type, "message_type_name": packet.message_type_name, "message_subtype": packet.message_subtype, "message_subtype_name": packet.message_subtype_name, "controller_timestamp_us": packet.controller_timestamp_us, "source": packet.source, "payload": packet.payload, "raw_hex": packet.raw.hex(), "strings_found": extract_printable_strings(packet.raw), "suspected_error_code": detect_ur_error_code(*extract_printable_strings(packet.raw)), }, ) def record_parse_error(self, raw: bytes, received_at: str, exc: Exception, stage: str = "parse_or_apply") -> None: self.write( "parse_error", { "stage": stage, "received_at": received_at, "error_type": type(exc).__name__, "error": str(exc), "traceback": traceback.format_exc(), "raw_packet": raw_packet_summary(raw, received_at), }, ) def record_message_entry(self, entry: dict[str, Any], raw_hex: str | None = None) -> None: if entry.get("detected_error_code") or entry.get("robot_message_code") not in (None, 0): payload = dict(entry) if raw_hex is not None: payload["raw_hex"] = raw_hex self.write("message_event", payload) def recent_packets(self) -> list[dict[str, Any]]: with self._lock: return deep_copy_jsonable(list(self._recent_packets)) @property def last_export_path(self) -> str: with self._lock: return self._last_export_path def export_bundle(self, client: "URClient", note: str = "") -> str: timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") zip_path = self.export_dir / f"ur_debug_bundle_{timestamp}.zip" snapshot = client.snapshot() messages = [item.to_dict() for item in client.messages(limit=None)] settings_dict = json_safe(client.settings.__dict__) recent_packets = self.recent_packets() with zipfile.ZipFile(zip_path, "w", compression=zipfile.ZIP_DEFLATED) as zf: zf.writestr("snapshot.json", json.dumps(snapshot, ensure_ascii=False, indent=2)) zf.writestr("messages.json", json.dumps(messages, ensure_ascii=False, indent=2)) zf.writestr("recent_packets.json", json.dumps(recent_packets, ensure_ascii=False, indent=2)) zf.writestr("settings.json", json.dumps(settings_dict, ensure_ascii=False, indent=2)) zf.writestr("note.txt", note) if self.path is not None and self.path.exists(): zf.write(self.path, arcname=self.path.name) for maybe_path in [client.store.ndjson_path, client.store.snapshot_json_path, client.store.sqlite_path]: if maybe_path is not None and Path(maybe_path).exists(): zf.write(maybe_path, arcname=Path(maybe_path).name) with self._lock: self._last_export_path = str(zip_path) self.write("bundle_exported", {"zip_path": str(zip_path), "note": note}) return str(zip_path) def _flatten(prefix: str, value: Any, out: dict[str, Any]) -> None: if isinstance(value, dict): if not value and prefix: out[prefix] = {} return for key in sorted(value.keys()): next_prefix = f"{prefix}.{key}" if prefix else str(key) _flatten(next_prefix, value[key], out) return if isinstance(value, list): if not value and prefix: out[prefix] = [] return if all(not isinstance(item, (dict, list)) for item in value): out[prefix] = value return for index, item in enumerate(value): next_prefix = f"{prefix}[{index}]" if prefix else f"[{index}]" _flatten(next_prefix, item, out) return if prefix: out[prefix] = value def flatten_dict(value: dict[str, Any]) -> dict[str, Any]: out: dict[str, Any] = {} _flatten("", value, out) return out def _joint_dict(values: list[Any], key: str | None = None) -> dict[str, Any]: result: dict[str, Any] = {} for index, label in enumerate(JOINT_LABELS): if index >= len(values): break if key is None: result[label] = values[index] else: result[label] = {key: values[index]} return result def _merge_joint_arrays(*dicts: dict[str, dict[str, Any]]) -> dict[str, dict[str, Any]]: merged: dict[str, dict[str, Any]] = {label: {} for label in JOINT_LABELS} for item in dicts: for label, payload in item.items(): if isinstance(payload, dict): merged.setdefault(label, {}).update(payload) else: merged.setdefault(label, {})["value"] = payload return {label: values for label, values in merged.items() if values} class URPacketParser: def __init__(self, context: ParseContext | None = None) -> None: self.context = context or ParseContext() def parse_packet(self, raw: bytes, received_at: str | None = None) -> ParsedPacket: if len(raw) < 5: raise PacketLengthError(f"packet too short: {len(raw)}") reader = ByteReader(raw) declared_size = reader.read_uint32() if declared_size != len(raw): raise PacketLengthError(f"declared packet size {declared_size} != received {len(raw)}") message_type = reader.read_uint8() packet = ParsedPacket( received_at=received_at or utc_now_iso(), message_type=message_type, message_type_name=enum_name(MESSAGE_TYPES, message_type, "unknown_message_type"), message_subtype=None, message_subtype_name=None, controller_timestamp_us=None, source=None, payload={}, raw=raw, ) if message_type == 16: packet.payload, packet.controller_timestamp_us = self._parse_robot_state(reader) elif message_type == 20: packet.message_subtype, packet.message_subtype_name, packet.controller_timestamp_us, packet.source, packet.payload = self._parse_robot_message(reader) elif message_type == 25: packet.message_subtype, packet.message_subtype_name, packet.controller_timestamp_us, packet.payload = self._parse_program_state_message(reader) else: packet.payload = {"unparsed": True, "raw_hex": reader.remaining_bytes().hex()} return packet @staticmethod def _read_optional(reader: ByteReader, minimum_size: int, fn: Any) -> Any | None: return fn() if reader.remaining() >= minimum_size else None @staticmethod def _read_pose64(reader: ByteReader, prefix: str = "") -> dict[str, float]: labels = ["x", "y", "z", "rx", "ry", "rz"] values = [reader.read_double() for _ in range(6)] return {f"{prefix}{label}": value for label, value in zip(labels, values)} @staticmethod def _read_pose32(reader: ByteReader) -> dict[str, float]: labels = ["x", "y", "z", "rx", "ry", "rz"] values = [reader.read_float() for _ in range(6)] return {label: value for label, value in zip(labels, values)} def _parse_robot_state(self, reader: ByteReader) -> tuple[dict[str, Any], int | None]: subpackages: dict[str, Any] = {} timestamp: int | None = None while not reader.empty(): if reader.remaining() < 5: raise ParseError("robot_state subpackage header too short") package_size = reader.read_uint32() package_type = reader.read_uint8() if package_size < 5: raise ParseError(f"invalid subpackage size: {package_size}") body_size = package_size - 5 if body_size > reader.remaining(): raise ParseError(f"subpackage exceeds remaining bytes: {package_size}") body = reader.read_reader(body_size) package_name = enum_name(ROBOT_STATE_SUBPACKAGES, package_type, "unknown_subpackage") payload = self._parse_robot_state_body(package_type, body) subpackages[package_name] = { "package_type": package_type, "package_type_name": package_name, "payload": payload, } if package_name == "robot_mode_data" and isinstance(payload.get("timestamp"), int): timestamp = payload["timestamp"] return {"subpackages": subpackages}, timestamp def _parse_robot_state_body(self, package_type: int, reader: ByteReader) -> dict[str, Any]: parsers = { 0: self._parse_state_robot_mode_data, 1: self._parse_state_joint_data, 2: self._parse_state_tool_data, 3: self._parse_state_masterboard_data, 4: self._parse_state_cartesian_info, 5: self._parse_state_kinematics_info, 6: self._parse_state_configuration_data, 7: self._parse_state_force_mode_data, 8: self._parse_state_additional_info, 9: self._parse_state_internal_skip, 10: self._parse_state_internal_skip, 11: self._parse_state_tool_communication_info, 12: self._parse_state_tool_mode_info, 13: self._parse_state_singularity_info, } parser = parsers.get(package_type) if parser is None: return {"unparsed": True, "raw_hex": reader.remaining_bytes().hex()} payload = parser(reader) if not reader.empty(): payload["trailing_bytes_hex"] = reader.remaining_bytes().hex() reader.skip(reader.remaining()) return payload def _parse_state_robot_mode_data(self, reader: ByteReader) -> dict[str, Any]: data = { "timestamp": reader.read_uint64(), "is_real_robot_connected": reader.read_bool(), "is_real_robot_enabled": reader.read_bool(), "is_robot_power_on": reader.read_bool(), "is_emergency_stopped": reader.read_bool(), "is_protective_stopped": reader.read_bool(), "is_program_running": reader.read_bool(), "is_program_paused": reader.read_bool(), } robot_mode = reader.read_uint8() control_mode = reader.read_uint8() data["robot_mode"] = robot_mode data["robot_mode_name"] = enum_name(ROBOT_MODES, robot_mode, "unknown_robot_mode") data["control_mode"] = control_mode data["control_mode_name"] = enum_name(CONTROL_MODES, control_mode, "unknown_control_mode") data["target_speed_fraction"] = reader.read_double() data["speed_scaling"] = reader.read_double() target_limit = self._read_optional(reader, 8, reader.read_double) if target_limit is not None: data["target_speed_fraction_limit"] = target_limit reserved = self._read_optional(reader, 1, reader.read_uint8) if reserved is not None: data["reserved"] = reserved return data def _parse_state_joint_data(self, reader: ByteReader) -> dict[str, Any]: joints: dict[str, Any] = {} bytes_per_joint = 41 for label in JOINT_LABELS: if reader.remaining() < bytes_per_joint: raise ParseError(f"joint_data incomplete for {label}") mode = None joints[label] = { "q_actual": reader.read_double(), "q_target": reader.read_double(), "qd_actual": reader.read_double(), "i_actual": reader.read_float(), "v_actual": reader.read_float(), "t_motor": reader.read_float(), "t_micro_deprecated": reader.read_float(), } mode = reader.read_uint8() joints[label]["joint_mode"] = mode joints[label]["joint_mode_name"] = enum_name(JOINT_MODES, mode, "unknown_joint_mode") return {"joints": joints} def _parse_state_tool_data(self, reader: ByteReader) -> dict[str, Any]: data = { "analog_input_range_0": reader.read_uint8(), "analog_input_range_1": reader.read_uint8(), "analog_input_0": reader.read_double(), "analog_input_1": reader.read_double(), "tool_voltage_48v": reader.read_float(), "tool_output_voltage": reader.read_uint8(), "tool_current": reader.read_float(), "tool_temperature": reader.read_float(), } tool_mode = reader.read_uint8() data["tool_mode"] = tool_mode data["tool_mode_name"] = enum_name(TOOL_MODES, tool_mode, "unknown_tool_mode") return data def _parse_state_masterboard_data(self, reader: ByteReader) -> dict[str, Any]: data = { "digital_input_bits": reader.read_uint32(), "digital_output_bits": reader.read_uint32(), "analog_input_range_0": reader.read_uint8(), "analog_input_range_1": reader.read_uint8(), "analog_input_0": reader.read_double(), "analog_input_1": reader.read_double(), "analog_output_domain_0": reader.read_int8(), "analog_output_domain_1": reader.read_int8(), "analog_output_0": reader.read_double(), "analog_output_1": reader.read_double(), "masterboard_temperature": reader.read_float(), "robot_voltage_48v": reader.read_float(), "robot_current": reader.read_float(), "master_io_current": reader.read_float(), } safety_mode = reader.read_uint8() data["safety_mode"] = safety_mode data["safety_mode_name"] = enum_name(SAFETY_MODES, safety_mode, "unknown_safety_mode") data["in_reduced_mode"] = bool(reader.read_uint8()) euromap_installed = bool(reader.read_int8()) data["euromap67_interface_installed"] = euromap_installed if euromap_installed and reader.remaining() >= 16: data["euromap67"] = { "input_bits": reader.read_uint32(), "output_bits": reader.read_uint32(), "voltage_24v": reader.read_float(), "current": reader.read_float(), } maybe_internal = self._read_optional(reader, 4, reader.read_uint32) if maybe_internal is not None: data["ur_internal_input_bitmask"] = maybe_internal selector = self._read_optional(reader, 1, reader.read_uint8) if selector is not None: data["operational_mode_selector_input"] = selector three_position = self._read_optional(reader, 1, reader.read_uint8) if three_position is not None: data["three_position_enabling_device_input"] = three_position internal_byte = self._read_optional(reader, 1, reader.read_uint8) if internal_byte is not None: data["ur_internal_byte"] = internal_byte return data def _parse_state_cartesian_info(self, reader: ByteReader) -> dict[str, Any]: data = {"tcp_pose": self._read_pose64(reader)} offset = self._read_optional(reader, 48, lambda: self._read_pose64(reader)) if offset is not None: data["tcp_offset"] = offset return data def _parse_state_kinematics_info(self, reader: ByteReader) -> dict[str, Any]: checksums = [reader.read_uint32() for _ in range(6)] dh_theta = [reader.read_double() for _ in range(6)] dh_a = [reader.read_double() for _ in range(6)] dh_d = [reader.read_double() for _ in range(6)] dh_alpha = [reader.read_double() for _ in range(6)] calibration_status = reader.read_uint32() return { "checksum": {label: checksums[index] for index, label in enumerate(JOINT_LABELS)}, "joints": _merge_joint_arrays( _joint_dict(dh_theta, "dh_theta"), _joint_dict(dh_a, "dh_a"), _joint_dict(dh_d, "dh_d"), _joint_dict(dh_alpha, "dh_alpha"), ), "calibration_status": calibration_status, } def _parse_state_configuration_data(self, reader: ByteReader) -> dict[str, Any]: joint_min = [reader.read_double() for _ in range(6)] joint_max = [reader.read_double() for _ in range(6)] joint_max_speed = [reader.read_double() for _ in range(6)] joint_max_acc = [reader.read_double() for _ in range(6)] data = { "joints": _merge_joint_arrays( _joint_dict(joint_min, "joint_min_limit"), _joint_dict(joint_max, "joint_max_limit"), _joint_dict(joint_max_speed, "joint_max_speed"), _joint_dict(joint_max_acc, "joint_max_acceleration"), ), "v_joint_default": reader.read_double(), "a_joint_default": reader.read_double(), "v_tool_default": reader.read_double(), "a_tool_default": reader.read_double(), "eq_radius": reader.read_double(), } dh_a = [reader.read_double() for _ in range(6)] dh_d = [reader.read_double() for _ in range(6)] dh_alpha = [reader.read_double() for _ in range(6)] dh_theta = [reader.read_double() for _ in range(6)] for label, values in _merge_joint_arrays( _joint_dict(dh_a, "dh_a"), _joint_dict(dh_d, "dh_d"), _joint_dict(dh_alpha, "dh_alpha"), _joint_dict(dh_theta, "dh_theta"), ).items(): data["joints"].setdefault(label, {}).update(values) data["masterboard_version"] = reader.read_int32() data["controller_box_type"] = reader.read_int32() data["robot_type"] = reader.read_int32() data["robot_sub_type"] = reader.read_int32() return data def _parse_state_force_mode_data(self, reader: ByteReader) -> dict[str, Any]: data = {"force_pose": self._read_pose64(reader, prefix="f_")} dexterity = self._read_optional(reader, 8, reader.read_double) if dexterity is not None: data["robot_dexterity"] = dexterity return data def _parse_state_additional_info(self, reader: ByteReader) -> dict[str, Any]: data = { "tp_button_state": reader.read_uint8(), "freedrive_button_enabled": reader.read_bool(), "io_enabled_freedrive": reader.read_bool(), } reserved = self._read_optional(reader, 1, reader.read_uint8) if reserved is not None: data["reserved"] = reserved return data def _parse_state_internal_skip(self, reader: ByteReader) -> dict[str, Any]: return {"skipped_internal_package": True, "raw_hex": reader.remaining_bytes().hex()} def _parse_state_tool_communication_info(self, reader: ByteReader) -> dict[str, Any]: return { "tool_communication_is_enabled": reader.read_bool(), "baud_rate": reader.read_int32(), "parity": reader.read_int32(), "stop_bits": reader.read_int32(), "rx_idle_chars": reader.read_float(), "tx_idle_chars": reader.read_float(), } def _parse_state_tool_mode_info(self, reader: ByteReader) -> dict[str, Any]: return { "output_mode": reader.read_uint8(), "digital_output_mode_0": reader.read_uint8(), "digital_output_mode_1": reader.read_uint8(), } def _parse_state_singularity_info(self, reader: ByteReader) -> dict[str, Any]: return { "singularity_severity": reader.read_uint8(), "singularity_type": reader.read_uint8(), } def _parse_robot_message(self, reader: ByteReader) -> tuple[int, str, int, int, dict[str, Any]]: timestamp = reader.read_uint64() source = reader.read_int8() subtype = reader.read_int8() subtype_name = enum_name(ROBOT_MESSAGE_TYPES, subtype, "unknown_robot_message_type") if subtype == 0: text_message = reader.read_string() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "text_text_message": text_message, "text": text_message, } elif subtype == 2: request_id = reader.read_uint32() requested_type = reader.read_uint32() warning = reader.read_bool() error = reader.read_bool() blocking = reader.read_bool() title_size = reader.read_uint8() popup_title = reader.read_string(title_size) popup_text = reader.read_string() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "request_id": request_id, "requested_type": requested_type, "requested_type_name": enum_name(REQUESTED_TYPES, requested_type, "unknown_requested_type"), "warning": warning, "error": error, "blocking": blocking, "popup_message_title_size": title_size, "popup_message_title": popup_title, "popup_text_message": popup_text, "title": popup_title, "text": popup_text, } elif subtype == 3: project_name_size = reader.read_int8() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "project_name": reader.read_string(project_name_size), "major_version": reader.read_uint8(), "minor_version": reader.read_uint8(), "bugfix_version": reader.read_int32(), "build_number": reader.read_int32(), "build_date": reader.read_string(), } elif subtype == 5: robot_message_code = reader.read_int32() robot_message_argument = reader.read_int32() safety_mode_type = reader.read_uint8() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "robot_message_code": robot_message_code, "robot_message_argument": robot_message_argument, "safety_mode_type": safety_mode_type, "safety_mode_type_name": enum_name(SAFETY_MODES, safety_mode_type, "unknown_safety_mode"), } report_data_type = self._read_optional(reader, 4, reader.read_uint32) report_data = self._read_optional(reader, 4, reader.read_uint32) if report_data_type is not None: payload["report_data_type"] = report_data_type if report_data is not None: payload["report_data"] = report_data elif subtype == 6: robot_message_code = reader.read_int32() robot_message_argument = reader.read_int32() report_level = reader.read_int32() data_type = reader.read_uint32() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "robot_message_code": robot_message_code, "robot_message_argument": robot_message_argument, "robot_message_report_level": report_level, "robot_message_report_level_name": enum_name(REPORT_LEVELS, report_level, "unknown_report_level"), "robot_message_data_type": data_type, "robot_message_data_type_name": enum_name(ROBOT_MESSAGE_DATA_TYPES, data_type, "unknown_robot_message_data_type"), } if data_type in (0, 1, 4) and reader.remaining() >= 4: data_value = reader.read_uint32() payload["robot_message_data"] = data_value if data_type == 4: payload["robot_message_data_hex"] = f"0x{data_value:08X}" elif data_type == 2 and reader.remaining() >= 4: payload["robot_message_data"] = reader.read_int32() elif data_type == 3 and reader.remaining() >= 4: payload["robot_message_data"] = reader.read_float() elif data_type == 5: text_length = reader.read_uint16() if reader.remaining() >= 2 else 0 if text_length and reader.remaining() >= text_length: text_message = reader.read_string(text_length) else: text_message = reader.read_string() payload["text_length"] = text_length payload["error_text_message"] = text_message payload["text"] = text_message if not reader.empty(): payload["trailing_bytes_hex"] = reader.remaining_bytes().hex() reader.skip(reader.remaining()) elif subtype == 7: robot_message_code = reader.read_int32() robot_message_argument = reader.read_int32() title_size = reader.read_uint8() robot_message_title = reader.read_string(title_size) key_text_message = reader.read_string() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "robot_message_code": robot_message_code, "robot_message_argument": robot_message_argument, "robot_message_title_size": title_size, "robot_message_title": robot_message_title, "key_text_message": key_text_message, "title": robot_message_title, "text": key_text_message, } elif subtype == 9: request_id = reader.read_uint32() requested_type = reader.read_uint32() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "request_id": request_id, "requested_type": requested_type, "requested_type_name": enum_name(REQUESTED_TYPES, requested_type, "unknown_requested_type"), "text": reader.read_string(), } elif subtype == 10: script_line_number = reader.read_int32() script_column_number = reader.read_int32() runtime_exception_text = reader.read_string() payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "script_line_number": script_line_number, "script_column_number": script_column_number, "runtime_exception_text_message": runtime_exception_text, "text": runtime_exception_text, } elif subtype == 14: threads: list[dict[str, Any]] = [] while not reader.empty(): label_id = reader.read_int32() label_length = reader.read_int32() label_name = reader.read_string(label_length) thread_length = reader.read_int32() thread_name = reader.read_string(thread_length) threads.append( { "label_id": label_id, "label_name": label_name, "thread_name": thread_name, } ) payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "threads": threads, } else: payload = { "timestamp": timestamp, "source": source, "source_name": enum_name(MESSAGE_SOURCES, source, "unknown_source"), "unparsed": True, "raw_hex": reader.remaining_bytes().hex(), } reader.skip(reader.remaining()) return subtype, subtype_name, timestamp, source, payload def _parse_program_state_message(self, reader: ByteReader) -> tuple[int, str, int, dict[str, Any]]: timestamp = reader.read_uint64() subtype = reader.read_int8() subtype_name = enum_name(PROGRAM_STATE_MESSAGE_TYPES, subtype, "unknown_program_state_message_type") if subtype == 0: payload = self._parse_program_state_global_variables_setup(reader) elif subtype == 1: payload = self._parse_program_state_global_variables_update(reader) elif subtype == 2: title_size = reader.read_uint8() payload = {"title": reader.read_string(title_size), "text": reader.read_string()} else: payload = {"unparsed": True, "raw_hex": reader.remaining_bytes().hex()} reader.skip(reader.remaining()) return subtype, subtype_name, timestamp, payload def _parse_program_state_global_variables_setup(self, reader: ByteReader) -> dict[str, Any]: start_index = reader.read_uint16() raw_names = reader.read_string() names = [item for item in raw_names.split("\n") if item] self.context.set_variable_names(start_index, names) return { "start_index": start_index, "variable_names": names, "known_variable_count": len(self.context.global_variable_names), } def _parse_program_state_global_variables_update(self, reader: ByteReader) -> dict[str, Any]: start_index = reader.read_uint16() values: list[dict[str, Any]] = [] while not reader.empty(): values.append(self._parse_program_value(reader, expect_separator=True)) self.context.set_variable_values(start_index, values) resolved: list[dict[str, Any]] = [] for offset, value in enumerate(values): index = start_index + offset name = self.context.global_variable_names[index] if index < len(self.context.global_variable_names) else None resolved.append({"index": index, "name": name, "value": value.get("value")}) return {"start_index": start_index, "variables": resolved} def _parse_program_value(self, reader: ByteReader, expect_separator: bool) -> dict[str, Any]: type_code = reader.read_uint8() type_name = enum_name(PROGRAM_VALUE_TYPES, type_code, "unknown_program_value_type") if type_code == 0: payload = {"type_code": type_code, "type_name": type_name, "value": None} elif type_code in (3, 4): length = reader.read_int16() payload = {"type_code": type_code, "type_name": type_name, "value": reader.read_string(length)} elif type_code == 5: length = reader.read_int16() payload = { "type_code": type_code, "type_name": type_name, "value": [self._parse_program_value(reader, expect_separator=False)["value"] for _ in range(length)], } elif type_code == 10: payload = {"type_code": type_code, "type_name": type_name, "value": self._read_pose32(reader)} elif type_code == 12: payload = {"type_code": type_code, "type_name": type_name, "value": reader.read_bool()} elif type_code == 13: payload = {"type_code": type_code, "type_name": type_name, "value": reader.read_float()} elif type_code == 14: payload = {"type_code": type_code, "type_name": type_name, "value": reader.read_int32()} elif type_code == 15: payload = {"type_code": type_code, "type_name": type_name, "value": reader.read_float()} else: raw = reader.remaining_bytes() payload = {"type_code": type_code, "type_name": type_name, "value": None, "raw_hex": raw.hex(), "unparsed": True} reader.skip(reader.remaining()) if expect_separator and not reader.empty(): separator = reader.read_uint8() if separator != 0x0A: raise ParseError(f"expected newline separator, got 0x{separator:02x}") return payload if expect_separator and not reader.empty(): separator = reader.read_uint8() if separator != 0x0A: raise ParseError(f"expected newline separator, got 0x{separator:02x}") return payload class Storage: def __init__(self, settings: AppSettings) -> None: self.settings = settings self.sqlite_path = ensure_parent(settings.sqlite_path) if settings.save_sqlite else None self.ndjson_path = ensure_parent(settings.ndjson_path) if settings.save_ndjson else None self.snapshot_json_path = ensure_parent(settings.snapshot_json_path) if settings.save_snapshot_json else None self._lock = threading.RLock() self._conn: sqlite3.Connection | None = None self._ndjson = None if self.sqlite_path is not None: self._conn = sqlite3.connect(self.sqlite_path, timeout=30.0, check_same_thread=False) self._conn.execute("PRAGMA journal_mode=WAL") self._conn.execute("PRAGMA synchronous=NORMAL") self._ensure_schema() if self.ndjson_path is not None: self._ndjson = self.ndjson_path.open("a", encoding="utf-8", buffering=1) def close(self) -> None: with self._lock: if self._ndjson is not None: self._ndjson.close() self._ndjson = None if self._conn is not None: self._conn.close() self._conn = None def reset_files(self) -> None: self.close() for path in [self.sqlite_path, self.ndjson_path, self.snapshot_json_path]: if path is not None and path.exists(): path.unlink() if self.settings.save_sqlite and self.sqlite_path is not None: self._conn = sqlite3.connect(self.sqlite_path, timeout=30.0, check_same_thread=False) self._conn.execute("PRAGMA journal_mode=WAL") self._conn.execute("PRAGMA synchronous=NORMAL") self._ensure_schema() if self.settings.save_ndjson and self.ndjson_path is not None: self._ndjson = self.ndjson_path.open("a", encoding="utf-8", buffering=1) def _table_exists(self, name: str) -> bool: assert self._conn is not None row = self._conn.execute( "SELECT name FROM sqlite_master WHERE type='table' AND name=?", (name,), ).fetchone() return row is not None def _table_columns(self, name: str) -> set[str]: assert self._conn is not None return {str(row[1]) for row in self._conn.execute(f"PRAGMA table_info({name})").fetchall()} def _recreate_if_incompatible(self, name: str, required_columns: set[str], ddl: str) -> None: assert self._conn is not None if not self._table_exists(name): self._conn.execute(ddl) return current = self._table_columns(name) if required_columns.issubset(current): return backup = f"{name}_legacy_{int(time.time())}" self._conn.execute(f"ALTER TABLE {name} RENAME TO {backup}") self._conn.execute(ddl) def _ensure_schema(self) -> None: assert self._conn is not None self._recreate_if_incompatible( "events", { "id", "received_at", "message_type", "message_type_name", "message_subtype", "message_subtype_name", "controller_timestamp_us", "source", "payload_json", "raw_hex", }, """ CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, message_type INTEGER NOT NULL, message_type_name TEXT NOT NULL, message_subtype INTEGER, message_subtype_name TEXT, controller_timestamp_us TEXT, source INTEGER, payload_json TEXT NOT NULL, raw_hex TEXT NOT NULL ) """, ) self._recreate_if_incompatible( "latest_snapshot", {"snapshot_key", "updated_at", "snapshot_json"}, """ CREATE TABLE IF NOT EXISTS latest_snapshot ( snapshot_key TEXT PRIMARY KEY, updated_at TEXT NOT NULL, snapshot_json TEXT NOT NULL ) """, ) self._recreate_if_incompatible( "latest_values", {"field_path", "updated_at", "value_json", "value_text"}, """ CREATE TABLE IF NOT EXISTS latest_values ( field_path TEXT PRIMARY KEY, updated_at TEXT NOT NULL, value_json TEXT NOT NULL, value_text TEXT NOT NULL ) """, ) self._conn.commit() def write_event(self, packet: ParsedPacket) -> None: with self._lock: if self._conn is not None: self._conn.execute( """ INSERT INTO events ( received_at, message_type, message_type_name, message_subtype, message_subtype_name, controller_timestamp_us, source, payload_json, raw_hex ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( packet.received_at, packet.message_type, packet.message_type_name, packet.message_subtype, packet.message_subtype_name, str(packet.controller_timestamp_us) if packet.controller_timestamp_us is not None else None, packet.source, json.dumps(json_safe(packet.payload), ensure_ascii=False), packet.raw.hex(), ), ) self._conn.commit() if self._ndjson is not None: self._ndjson.write(json.dumps(packet.event_dict(), ensure_ascii=False) + "\n") self._ndjson.flush() def write_snapshot(self, snapshot: dict[str, Any], flat_values: dict[str, Any]) -> None: updated_at = utc_now_iso() snapshot_json = json.dumps(json_safe(snapshot), ensure_ascii=False) with self._lock: if self._conn is not None: self._conn.execute( "INSERT OR REPLACE INTO latest_snapshot(snapshot_key, updated_at, snapshot_json) VALUES (?, ?, ?)", ("main", updated_at, snapshot_json), ) self._conn.execute("DELETE FROM latest_values") rows = [ ( field_path, updated_at, json.dumps(json_safe(value), ensure_ascii=False), format_value(value), ) for field_path, value in sorted(flat_values.items()) ] if rows: self._conn.executemany( "INSERT OR REPLACE INTO latest_values(field_path, updated_at, value_json, value_text) VALUES (?, ?, ?, ?)", rows, ) self._conn.commit() if self.snapshot_json_path is not None: self.snapshot_json_path.write_text(snapshot_json, encoding="utf-8")
[docs] class URClient: def __init__(self, settings: AppSettings | None = None, **settings_overrides: Any) -> None: if settings is None: settings = AppSettings(**settings_overrides) elif settings_overrides: merged = {**settings.__dict__, **settings_overrides} settings = AppSettings(**merged) self.settings = settings self.port = resolve_port(settings.interface_name, settings.port) self.parser = URPacketParser() self.store = Storage(settings) self.diagnostics = Diagnostics(settings) self.error_db = ErrorCodeDirectory(settings.error_db_path if settings.error_db_enabled else None) self._lock = threading.RLock() self._thread: threading.Thread | None = None self._stop_event = threading.Event() self._update_event = threading.Event() self._update_counter = 0 self._callbacks: list[Any] = [] self._socket: socket.socket | None = None self._messages = deque(maxlen=max(1, settings.message_history_limit)) self._snapshot: dict[str, Any] = { "meta": { "running": False, "connected": False, "status": "stopped", "host": settings.host, "port": self.port, "interface_name": settings.interface_name, "interface_label": interface_label(settings.interface_name), "started_at": "", "last_heartbeat_at": "", "last_packet_at": "", "packet_count": 0, "last_error": "", "diagnostics_enabled": bool(settings.diagnostics_enabled), "diagnostics_path": str(Path(settings.diagnostics_path).expanduser()), "last_debug_export": "", }, "values": { "robot": {}, "tcp": {}, "joints": {}, "masterboard": {}, "tool": {}, "kinematics": {}, "configuration": {}, "additional": {}, "tool_communication": {}, "tool_mode": {}, "singularity": {}, "globals": {}, "version": {}, "messages": {}, "connection": { "supports_key_messages": supports_key_messages(settings.interface_name), "interface_note": interface_note(settings.interface_name), }, "error_db": { "enabled": bool(settings.error_db_enabled), "path": "" if self.error_db.path is None else str(self.error_db.path), "available": bool(self.error_db.available), "last_error": self.error_db.last_error, "last_lookup_code": "", "last_lookup_found": False, }, "errors": {"latest": {}, "latest_message": {}, "latest_raw": {}}, "diagnostics": {"path": str(Path(settings.diagnostics_path).expanduser()), "last_export_path": "", "recent_packet_count": 0, "latest_raw_detected_error": {}}, "packets": {"last": {}}, }, "messages": [], } self._flat_values: dict[str, Any] = {} self._last_state_write_monotonic = 0.0 self._refresh_flat_values() @property def running(self) -> bool: with self._lock: return bool(self._snapshot["meta"]["running"]) @property def connected(self) -> bool: with self._lock: return bool(self._snapshot["meta"]["connected"]) @property def packet_count(self) -> int: with self._lock: return int(self._snapshot["meta"]["packet_count"]) @property def status(self) -> str: with self._lock: return str(self._snapshot["meta"]["status"])
[docs] def read(self, path: str, default: Any = None) -> Any: return self.get(path, default)
[docs] def read_many(self, *paths: str) -> dict[str, Any]: return {path: self.get(path, None) for path in paths}
[docs] def select(self, *paths: str) -> dict[str, Any]: return self.read_many(*paths)
def __getitem__(self, path: str) -> Any: value = self.get(path, None) if value is None and path not in self.keys(): raise KeyError(path) return value
[docs] def wait_until_connected(self, timeout: float = 5.0) -> bool: deadline = time.monotonic() + max(0.0, timeout) while time.monotonic() <= deadline: if self.connected: return True time.sleep(0.05) return self.connected
[docs] def wait_for_packet(self, timeout: float = 5.0) -> bool: return self.wait_for_update(timeout=timeout) is not None
[docs] def wait_for_packets(self, minimum_packet_count: int = 1, timeout: float = 5.0) -> bool: deadline = time.monotonic() + max(0.0, timeout) while time.monotonic() <= deadline: if self.packet_count >= minimum_packet_count: return True time.sleep(0.05) return self.packet_count >= minimum_packet_count
[docs] def subscribe(self, callback: Any) -> None: with self._lock: if callback not in self._callbacks: self._callbacks.append(callback)
[docs] def unsubscribe(self, callback: Any) -> None: with self._lock: if callback in self._callbacks: self._callbacks.remove(callback)
[docs] def wait_for_update(self, last_packet_count: int | None = None, timeout: float = 5.0) -> int | None: baseline = self.packet_count if last_packet_count is None else int(last_packet_count) deadline = time.monotonic() + max(0.0, timeout) while time.monotonic() <= deadline: if self.packet_count > baseline: return self.packet_count time.sleep(0.05) return None
@property def robot_mode(self) -> Any: return self.read("robot.mode") @property def safety_mode(self) -> Any: return self.read("robot.safety_mode") @property def tcp_pose(self) -> dict[str, Any]: pose = self.read("tcp.pose", {}) return pose if isinstance(pose, dict) else {} @property def joint_positions(self) -> dict[str, Any]: result: dict[str, Any] = {} for label in JOINT_LABELS: value = self.read(f"joints.{label}.q_actual") if value is not None: result[label] = value return result @property def global_variables(self) -> dict[str, Any]: values = self.read("globals", {}) return values if isinstance(values, dict) else {}
[docs] def get_variable(self, name: str, default: Any = None) -> Any: return self.read(f"globals.{name}", default)
[docs] def watch(self, *paths: str, interval: float = 0.5) -> Iterator[dict[str, Any]]: while self.running or self.connected: self.wait_for_update(timeout=max(0.0, interval)) if paths: yield self.read_many(*paths) else: yield self.snapshot()
[docs] def messages(self, category: str | None = None, limit: int | None = None) -> list[MessageRecord]: with self._lock: raw_messages = list(self._messages) records = [MessageRecord(**deep_copy_jsonable(entry)) for entry in raw_messages if category is None or entry.get("category") == category] if limit is not None: return records[: max(0, int(limit))] return records
[docs] def latest_message(self, category: str | None = None) -> MessageRecord | None: items = self.messages(category=category, limit=1) return items[0] if items else None
@property def latest_key_message(self) -> KeyMessage | None: message = self.latest_message("key") if message is None: return None payload = message.payload return KeyMessage( received_at=message.received_at, timestamp=payload.get("timestamp"), source=payload.get("source"), source_name=payload.get("source_name"), robot_message_code=payload.get("robot_message_code"), robot_message_argument=payload.get("robot_message_argument"), robot_message_title_size=payload.get("robot_message_title_size"), robot_message_title=str(payload.get("robot_message_title") or payload.get("title") or ""), key_text_message=str(payload.get("key_text_message") or payload.get("text") or ""), )
[docs] def latest_key_message_dict(self) -> dict[str, Any] | None: message = self.latest_key_message return None if message is None else message.to_dict()
@property def latest_code_message(self) -> MessageRecord | None: with self._lock: entry = self._snapshot.get("values", {}).get("messages", {}).get("latest_with_code") return None if not entry else MessageRecord(**deep_copy_jsonable(entry))
[docs] def latest_code_message_dict(self) -> dict[str, Any] | None: message = self.latest_code_message return None if message is None else message.to_dict()
@property def latest_nonzero_code_message(self) -> MessageRecord | None: with self._lock: entry = self._snapshot.get("values", {}).get("messages", {}).get("latest_nonzero_code") return None if not entry else MessageRecord(**deep_copy_jsonable(entry))
[docs] def latest_nonzero_code_message_dict(self) -> dict[str, Any] | None: message = self.latest_nonzero_code_message return None if message is None else message.to_dict()
@property def latest_detected_error_message(self) -> MessageRecord | None: with self._lock: entry = self._snapshot.get("values", {}).get("messages", {}).get("latest_detected_error") return None if not entry else MessageRecord(**deep_copy_jsonable(entry))
[docs] def latest_detected_error_message_dict(self) -> dict[str, Any] | None: message = self.latest_detected_error_message return None if message is None else message.to_dict()
@property def latest_popup_message(self) -> MessageRecord | None: with self._lock: entry = self._snapshot.get("values", {}).get("messages", {}).get("latest_popup") return None if not entry else MessageRecord(**deep_copy_jsonable(entry))
[docs] def latest_popup_message_dict(self) -> dict[str, Any] | None: message = self.latest_popup_message return None if message is None else message.to_dict()
[docs] def lookup_error(self, code: str | None) -> dict[str, Any] | None: return self.error_db.lookup(code)
[docs] def search_errors(self, query: str, limit: int = 20) -> list[dict[str, Any]]: return self.error_db.search(query, limit=limit)
[docs] def latest_error_details(self) -> dict[str, Any] | None: with self._lock: entry = self._snapshot.get("values", {}).get("errors", {}).get("latest") return None if not entry else deep_copy_jsonable(entry)
[docs] def latest_message_error_details(self) -> dict[str, Any] | None: with self._lock: entry = self._snapshot.get("values", {}).get("errors", {}).get("latest_message") return None if not entry else deep_copy_jsonable(entry)
[docs] def latest_raw_error_details(self) -> dict[str, Any] | None: with self._lock: entry = self._snapshot.get("values", {}).get("errors", {}).get("latest_raw") return None if not entry else deep_copy_jsonable(entry)
[docs] def messages_with_code(self, nonzero_only: bool = False, limit: int | None = None) -> list[MessageRecord]: records = [item for item in self.messages(limit=None) if item.robot_message_code is not None] if nonzero_only: records = [item for item in records if item.robot_message_code not in (None, 0)] if limit is not None: return records[: max(0, int(limit))] return records
[docs] def keys(self) -> list[str]: with self._lock: return sorted(self._flat_values.keys())
[docs] def flat_values(self) -> dict[str, Any]: with self._lock: return deep_copy_jsonable(self._flat_values)
[docs] def values(self) -> dict[str, Any]: with self._lock: return deep_copy_jsonable(self._snapshot["values"])
[docs] def meta(self) -> dict[str, Any]: with self._lock: return deep_copy_jsonable(self._snapshot["meta"])
[docs] def get(self, path: str, default: Any = None) -> Any: with self._lock: if path in self._flat_values: return deep_copy_jsonable(self._flat_values[path]) return deep_copy_jsonable(default)
[docs] def snapshot(self) -> dict[str, Any]: with self._lock: return deep_copy_jsonable(self._snapshot)
[docs] def start(self) -> None: if self._thread is not None and self._thread.is_alive(): return self.port = resolve_port(self.settings.interface_name, self.settings.port) with self._lock: self._snapshot["meta"].update( { "running": True, "connected": False, "status": "starting", "host": self.settings.host, "port": self.port, "interface_name": self.settings.interface_name, "interface_label": interface_label(self.settings.interface_name), "started_at": utc_now_iso(), "last_error": "", "last_heartbeat_at": utc_now_iso(), } ) self._snapshot["values"].setdefault("connection", {})["supports_key_messages"] = supports_key_messages(self.settings.interface_name) self._snapshot["values"].setdefault("connection", {})["interface_note"] = interface_note(self.settings.interface_name) self._snapshot["values"].setdefault("diagnostics", {})["path"] = str(Path(self.settings.diagnostics_path).expanduser()) self._snapshot["values"].setdefault("diagnostics", {})["recent_packet_count"] = len(self.diagnostics.recent_packets()) self._snapshot["values"].setdefault("diagnostics", {})["last_export_path"] = self.diagnostics.last_export_path self._refresh_flat_values() self._flush_state(force=True) self._stop_event.clear() self._thread = threading.Thread(target=self._worker, name="URClientMonitor", daemon=True) self._thread.start()
[docs] def stop(self) -> None: self._stop_event.set() if self._socket is not None: try: self._socket.shutdown(socket.SHUT_RDWR) except OSError: pass try: self._socket.close() except OSError: pass if self._thread is not None and self._thread.is_alive(): self._thread.join(timeout=3.0) with self._lock: self._snapshot["meta"].update( { "running": False, "connected": False, "status": "stopped", "last_heartbeat_at": utc_now_iso(), } ) self._flush_state(force=True)
[docs] def close(self) -> None: self.stop() self.store.close() self.diagnostics.close() self.error_db.close()
def __enter__(self) -> "URClient": self.start() return self def __exit__(self, exc_type: Any, exc: Any, tb: Any) -> None: self.close()
[docs] def reset_storage(self) -> None: self.stop() self.store.reset_files() self.diagnostics.reset_files() with self._lock: self._values().setdefault("diagnostics", {})["last_export_path"] = "" self._snapshot["meta"]["last_debug_export"] = "" self._flush_state(force=True)
[docs] def export_debug_bundle(self, note: str = "") -> str: path = self.diagnostics.export_bundle(self, note=note) with self._lock: self._snapshot["meta"]["last_debug_export"] = path self._values().setdefault("diagnostics", {})["last_export_path"] = path self._values().setdefault("diagnostics", {})["recent_packet_count"] = len(self.diagnostics.recent_packets()) self._flush_state(force=True) return path
def _build_error_details_entry(self, code: Any, source_kind: str, context: dict[str, Any]) -> dict[str, Any]: lookup = self.error_db.lookup(None if code is None else str(code)) normalized_code = normalize_ur_error_code(code) if lookup is None: lookup = { "input_code": "" if code is None else str(code), "normalized_code": normalized_code, "lookup_code": normalized_code, "found": False, "db_available": bool(self.error_db.available), "db_path": "" if self.error_db.path is None else str(self.error_db.path), "db_error": self.error_db.last_error or "", } code_value = str(lookup.get("error_code") or lookup.get("code") or lookup.get("lookup_code") or normalized_code or code or "") group_title = str(lookup.get("group_title") or "").strip() description = str(lookup.get("description") or "").strip() explanation = str(lookup.get("explanation") or "").strip() suggestion = str(lookup.get("suggestion") or "").strip() display_title = " | ".join(part for part in [code_value, group_title or description] if part) display_parts = [part for part in [description, explanation, suggestion] if part] entry = { "observed_at": utc_now_iso(), "source_kind": source_kind, "error_code": code_value, "code": code_value, "input_code": lookup.get("input_code"), "normalized_code": lookup.get("normalized_code"), "lookup_code": lookup.get("lookup_code"), "found": bool(lookup.get("found")), "group_code": lookup.get("group_code"), "group_title": lookup.get("group_title"), "description": lookup.get("description"), "explanation": lookup.get("explanation"), "suggestion": lookup.get("suggestion"), "page_start": lookup.get("page_start"), "page_end": lookup.get("page_end"), "section_no": lookup.get("section_no"), "source_document": lookup.get("source_document"), "db_available": bool(lookup.get("db_available")), "db_path": lookup.get("db_path"), "db_error": lookup.get("db_error"), "display_title": display_title or code_value, "display_text": "\n\n".join(display_parts), "context": deep_copy_jsonable(context), } return deep_copy_jsonable(entry) def _apply_error_lookup_locked(self, code: Any, source_kind: str, context: dict[str, Any]) -> None: if code in (None, ""): return entry = self._build_error_details_entry(code, source_kind=source_kind, context=context) values = self._values() values.setdefault("error_db", {}).update( { "enabled": bool(self.settings.error_db_enabled), "path": "" if self.error_db.path is None else str(self.error_db.path), "available": bool(self.error_db.available), "last_error": self.error_db.last_error, "last_lookup_code": entry.get("code") or entry.get("lookup_code") or entry.get("normalized_code") or "", "last_lookup_found": bool(entry.get("found")), } ) errors = values.setdefault("errors", {}) errors["latest"] = deep_copy_jsonable(entry) if source_kind == "message": errors["latest_message"] = deep_copy_jsonable(entry) elif source_kind == "raw_packet": errors["latest_raw"] = deep_copy_jsonable(entry) def _set_meta(self, **updates: Any) -> None: with self._lock: self._snapshot["meta"].update(updates) def _refresh_flat_values(self) -> None: self._flat_values = flatten_dict(self._snapshot.get("values", {})) meta_paths = { "meta.running": self._snapshot["meta"].get("running"), "meta.connected": self._snapshot["meta"].get("connected"), "meta.status": self._snapshot["meta"].get("status"), "meta.packet_count": self._snapshot["meta"].get("packet_count"), "meta.last_packet_at": self._snapshot["meta"].get("last_packet_at"), "meta.last_error": self._snapshot["meta"].get("last_error"), "meta.last_debug_export": self._snapshot["meta"].get("last_debug_export"), "meta.diagnostics_path": self._snapshot["meta"].get("diagnostics_path"), } self._flat_values.update(meta_paths) def _flush_state(self, force: bool = False) -> None: now = time.monotonic() if not force and (now - self._last_state_write_monotonic) < self.settings.state_write_interval: return with self._lock: self._refresh_flat_values() snapshot = deep_copy_jsonable(self._snapshot) flat_values = deep_copy_jsonable(self._flat_values) self.store.write_snapshot(snapshot, flat_values) self._last_state_write_monotonic = now def _worker(self) -> None: while not self._stop_event.is_set(): self._set_meta(status="connecting", connected=False, last_heartbeat_at=utc_now_iso()) self._flush_state(force=True) try: sock = socket.create_connection((self.settings.host, self.port), timeout=self.settings.connect_timeout) sock.settimeout(self.settings.read_timeout) sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1) self._socket = sock self._set_meta(status="connected", connected=True, last_error="", last_heartbeat_at=utc_now_iso()) self._flush_state(force=True) while not self._stop_event.is_set(): raw = self._read_packet(sock) received_at = utc_now_iso() raw_summary = self.diagnostics.record_raw_packet(raw, received_at) if raw_summary.get("suspected_error_code"): with self._lock: self._values().setdefault("diagnostics", {})["latest_raw_detected_error"] = deep_copy_jsonable(raw_summary) self._apply_error_lookup_locked(raw_summary.get("suspected_error_code"), source_kind="raw_packet", context=raw_summary) try: packet = self.parser.parse_packet(raw, received_at=received_at) self.diagnostics.record_parsed_packet(packet) self.store.write_event(packet) self._apply_packet(packet) except Exception as packet_exc: self.diagnostics.record_parse_error(raw, received_at, packet_exc) self._set_meta(last_error=f"packet parse/apply error: {packet_exc}", last_heartbeat_at=utc_now_iso()) with self._lock: self._values().setdefault("diagnostics", {})["recent_packet_count"] = len(self.diagnostics.recent_packets()) self._values().setdefault("diagnostics", {})["last_parse_error"] = {"received_at": received_at, "error": str(packet_exc), "error_type": type(packet_exc).__name__} self._flush_state(force=True) continue except Exception as exc: if self._stop_event.is_set(): break self._set_meta(status="error", connected=False, last_error=str(exc), last_heartbeat_at=utc_now_iso()) self._flush_state(force=True) time.sleep(self.settings.reconnect_delay) finally: try: if self._socket is not None: self._socket.close() except OSError: pass self._socket = None if not self._stop_event.is_set(): self._set_meta(status="disconnected", connected=False, last_heartbeat_at=utc_now_iso()) self._flush_state(force=True) self._set_meta(status="stopped", connected=False, running=False, last_heartbeat_at=utc_now_iso()) self._flush_state(force=True) def _recv_exactly(self, sock: socket.socket, size: int) -> bytes: data = bytearray() while len(data) < size: chunk = sock.recv(size - len(data)) if not chunk: raise OSError("socket closed by peer") data.extend(chunk) return bytes(data) def _read_packet(self, sock: socket.socket) -> bytes: header = self._recv_exactly(sock, 4) packet_length = struct.unpack(">I", header)[0] if packet_length < 5: raise PacketLengthError(f"invalid packet length: {packet_length}") if packet_length > self.settings.max_packet_size: raise PacketLengthError(f"packet length {packet_length} exceeds max {self.settings.max_packet_size}") body = self._recv_exactly(sock, packet_length - 4) return header + body def _apply_packet(self, packet: ParsedPacket) -> None: callbacks: list[Any] packet_event = packet.event_dict() with self._lock: self._snapshot["meta"]["packet_count"] = int(self._snapshot["meta"]["packet_count"]) + 1 self._snapshot["meta"]["last_packet_at"] = packet.received_at self._snapshot["meta"]["last_heartbeat_at"] = packet.received_at self._values().setdefault("packets", {})["last"] = deep_copy_jsonable(packet_event) if packet.message_type_name == "robot_state": self._apply_robot_state(packet.payload.get("subpackages", {})) elif packet.message_type_name == "robot_message": self._apply_robot_message(packet) elif packet.message_type_name == "program_state_message": self._apply_program_state_message(packet) self._values().setdefault("diagnostics", {})["recent_packet_count"] = len(self.diagnostics.recent_packets()) self._values().setdefault("diagnostics", {})["last_export_path"] = self.diagnostics.last_export_path self._refresh_flat_values() self._update_counter += 1 callbacks = list(self._callbacks) self._update_event.set() for callback in callbacks: try: callback(self, deep_copy_jsonable(packet_event)) except Exception: pass self._flush_state(force=False) def _values(self) -> dict[str, Any]: return self._snapshot["values"] def _apply_robot_state(self, subpackages: dict[str, Any]) -> None: values = self._values() robot_mode = subpackages.get("robot_mode_data", {}).get("payload", {}) if robot_mode: values["robot"].update( { "mode": robot_mode.get("robot_mode_name"), "control_mode": robot_mode.get("control_mode_name"), "safety_mode": values["robot"].get("safety_mode"), "target_speed_fraction": robot_mode.get("target_speed_fraction"), "speed_scaling": robot_mode.get("speed_scaling"), "target_speed_fraction_limit": robot_mode.get("target_speed_fraction_limit"), "is_real_robot_connected": robot_mode.get("is_real_robot_connected"), "is_real_robot_enabled": robot_mode.get("is_real_robot_enabled"), "is_robot_power_on": robot_mode.get("is_robot_power_on"), "is_emergency_stopped": robot_mode.get("is_emergency_stopped"), "is_protective_stopped": robot_mode.get("is_protective_stopped"), "is_program_running": robot_mode.get("is_program_running"), "is_program_paused": robot_mode.get("is_program_paused"), "timestamp": robot_mode.get("timestamp"), } ) joints = subpackages.get("joint_data", {}).get("payload", {}).get("joints", {}) if joints: values["joints"] = deep_copy_jsonable(joints) cartesian = subpackages.get("cartesian_info", {}).get("payload", {}) if cartesian: values["tcp"] = deep_copy_jsonable( { "pose": cartesian.get("tcp_pose", values["tcp"].get("pose", {})), "offset": cartesian.get("tcp_offset", values["tcp"].get("offset", {})), "force": values["tcp"].get("force", {}), "robot_dexterity": values["tcp"].get("robot_dexterity"), } ) masterboard = subpackages.get("masterboard_data", {}).get("payload", {}) if masterboard: values["masterboard"] = deep_copy_jsonable(masterboard) values["robot"]["safety_mode"] = masterboard.get("safety_mode_name") tool = subpackages.get("tool_data", {}).get("payload", {}) if tool: values["tool"] = deep_copy_jsonable(tool) kinematics = subpackages.get("kinematics_info", {}).get("payload", {}) if kinematics: values["kinematics"] = deep_copy_jsonable(kinematics) configuration = subpackages.get("configuration_data", {}).get("payload", {}) if configuration: values["configuration"] = deep_copy_jsonable(configuration) force = subpackages.get("force_mode_data", {}).get("payload", {}) if force: values.setdefault("tcp", {}).update( { "force": force.get("force_pose", {}), "robot_dexterity": force.get("robot_dexterity"), } ) additional = subpackages.get("additional_info", {}).get("payload", {}) if additional: values["additional"] = deep_copy_jsonable(additional) tool_comm = subpackages.get("tool_communication_info", {}).get("payload", {}) if tool_comm: values["tool_communication"] = deep_copy_jsonable(tool_comm) tool_mode = subpackages.get("tool_mode_info", {}).get("payload", {}) if tool_mode: values["tool_mode"] = deep_copy_jsonable(tool_mode) singularity = subpackages.get("singularity_info", {}).get("payload", {}) if singularity: values["singularity"] = deep_copy_jsonable(singularity) def _push_message(self, packet: ParsedPacket, category: str, title: str, text: str, payload: dict[str, Any]) -> None: formatted_error_code = format_robot_error_code(payload.get("robot_message_code"), payload.get("robot_message_argument")) detected_error_code = formatted_error_code or detect_ur_error_code( title, text, payload.get("robot_message_title"), payload.get("key_text_message"), payload.get("popup_message_title"), payload.get("popup_text_message"), payload.get("error_text_message"), payload.get("runtime_exception_text_message"), payload.get("text_text_message"), ) entry = { "received_at": packet.received_at, "category": category, "title": title, "text": text, "message_type_name": packet.message_type_name, "message_subtype_name": packet.message_subtype_name, "controller_timestamp_us": packet.controller_timestamp_us, "source": payload.get("source"), "source_name": payload.get("source_name"), "robot_message_code": payload.get("robot_message_code"), "robot_message_argument": payload.get("robot_message_argument"), "formatted_error_code": formatted_error_code, "detected_error_code": detected_error_code, "request_id": payload.get("request_id"), "requested_type": payload.get("requested_type_name") or payload.get("requested_type"), "payload": deep_copy_jsonable(payload), } self._messages.appendleft(entry) self._snapshot["messages"] = list(self._messages) values = self._values() message_values = values.setdefault("messages", {}) message_values["latest"] = deep_copy_jsonable(entry) by_category = message_values.setdefault("latest_by_category", {}) by_category[category] = deep_copy_jsonable(entry) if entry.get("robot_message_code") is not None: message_values["latest_with_code"] = deep_copy_jsonable(entry) if entry.get("robot_message_code") not in (0, None): message_values["latest_nonzero_code"] = deep_copy_jsonable(entry) if entry.get("detected_error_code"): message_values["latest_detected_error"] = deep_copy_jsonable(entry) if category == "popup": message_values["latest_popup"] = deep_copy_jsonable(entry) if category == "key": message_values["latest_key"] = { "received_at": packet.received_at, "timestamp": payload.get("timestamp"), "source": payload.get("source"), "source_name": payload.get("source_name"), "robot_message_code": payload.get("robot_message_code"), "robot_message_argument": payload.get("robot_message_argument"), "robot_message_title_size": payload.get("robot_message_title_size"), "robot_message_title": payload.get("robot_message_title") or payload.get("title") or "", "key_text_message": payload.get("key_text_message") or payload.get("text") or "", } if detected_error_code: self._apply_error_lookup_locked(detected_error_code, source_kind="message", context=entry) self.diagnostics.record_message_entry(entry, raw_hex=packet.raw.hex()) def _apply_robot_message(self, packet: ParsedPacket) -> None: values = self._values() payload = packet.payload subtype = packet.message_subtype_name or "robot_message" if subtype == "version": values["version"] = deep_copy_jsonable(payload) title = f"Version {payload.get('major_version')}.{payload.get('minor_version')}" text = f"{payload.get('project_name', '')} {payload.get('build_date', '')}".strip() self._push_message(packet, subtype, title, text, payload) return if payload.get("safety_mode_type_name"): values["robot"]["safety_mode"] = payload.get("safety_mode_type_name") if subtype == "program_threads": values["program_threads"] = deep_copy_jsonable(payload.get("threads", [])) title = str(payload.get("title") or subtype).strip() text = str(payload.get("text") or payload.get("build_date") or "").strip() formatted_error_code = format_robot_error_code(payload.get("robot_message_code"), payload.get("robot_message_argument")) if not text and formatted_error_code: text = formatted_error_code if not text and payload.get("robot_message_code") is not None: text = f"code={payload.get('robot_message_code')} arg={payload.get('robot_message_argument')}" if title or text: self._push_message(packet, subtype, title, text, payload) def _apply_program_state_message(self, packet: ParsedPacket) -> None: values = self._values() if packet.message_subtype_name == "global_variables_update": values["globals"] = deep_copy_jsonable(self.parser.context.named_variable_map()) self._push_message(packet, "global_variables_update", "Global variables updated", "", packet.payload) elif packet.message_subtype_name == "global_variables_setup": names = packet.payload.get("variable_names", []) self._push_message(packet, "global_variables_setup", "Global variables setup", ", ".join(names), packet.payload) elif packet.message_subtype_name == "variable_update": title = str(packet.payload.get("title", "Variable update")) text = str(packet.payload.get("text", "")) self._push_message(packet, "variable_update", title, text, packet.payload)
[docs] class ClientInterfaceAPI: """Simple class-based API for scripts and applications. Pass the robot IP address directly to the constructor, then call :meth:`read`, :meth:`read_many`, or use the convenience properties. """ def __init__( self, host: str, interface_name: str = "primary_ro", *, port: int | None = None, auto_start: bool = True, wait_ready: bool = False, wait_timeout: float = 5.0, minimum_packet_count: int = 1, error_db_path: str = "ur_error_codes.sqlite3", **settings_overrides: Any, ) -> None: merged = { "host": host, "interface_name": interface_name, "port": port, "error_db_path": error_db_path, **settings_overrides, } merged["auto_start"] = False self.settings = AppSettings(**merged) self.client = URClient(self.settings) if auto_start: self.start() if wait_ready: self.wait_until_ready(timeout=wait_timeout, minimum_packet_count=minimum_packet_count)
[docs] @classmethod def connect( cls, host: str, interface_name: str = "primary_ro", *, port: int | None = None, wait_ready: bool = True, wait_timeout: float = 5.0, minimum_packet_count: int = 1, **settings_overrides: Any, ) -> "ClientInterfaceAPI": return cls( host, interface_name=interface_name, port=port, auto_start=True, wait_ready=wait_ready, wait_timeout=wait_timeout, minimum_packet_count=minimum_packet_count, **settings_overrides, )
[docs] def start(self) -> "ClientInterfaceAPI": self.client.start() return self
[docs] def stop(self) -> None: self.client.stop()
[docs] def close(self) -> None: self.client.close()
def __enter__(self) -> "ClientInterfaceAPI": if not self.client.running: self.start() return self def __exit__(self, exc_type: Any, exc: Any, tb: Any) -> None: self.close() def __getitem__(self, path: str) -> Any: return self.client[path] @property def running(self) -> bool: return self.client.running @property def connected(self) -> bool: return self.client.connected @property def packet_count(self) -> int: return self.client.packet_count @property def status(self) -> str: return self.client.status @property def robot_mode(self) -> Any: return self.client.robot_mode @property def safety_mode(self) -> Any: return self.client.safety_mode @property def tcp_pose(self) -> dict[str, Any]: return self.client.tcp_pose @property def joint_positions(self) -> dict[str, Any]: return self.client.joint_positions @property def global_variables(self) -> dict[str, Any]: return self.client.global_variables @property def latest_key_message(self) -> KeyMessage | None: return self.client.latest_key_message @property def latest_error(self) -> dict[str, Any] | None: return self.client.latest_error_details()
[docs] def wait_until_ready(self, timeout: float = 5.0, minimum_packet_count: int = 1) -> bool: if not self.client.wait_until_connected(timeout): return False return self.client.wait_for_packets(minimum_packet_count=minimum_packet_count, timeout=timeout)
[docs] def wait_until_connected(self, timeout: float = 5.0) -> bool: return self.client.wait_until_connected(timeout)
[docs] def wait_for_packet(self, timeout: float = 5.0) -> bool: return self.client.wait_for_packet(timeout)
[docs] def wait_for_packets(self, minimum_packet_count: int = 1, timeout: float = 5.0) -> bool: return self.client.wait_for_packets(minimum_packet_count=minimum_packet_count, timeout=timeout)
[docs] def wait_for_update(self, last_packet_count: int | None = None, timeout: float = 5.0) -> int | None: return self.client.wait_for_update(last_packet_count=last_packet_count, timeout=timeout)
[docs] def read(self, path: str, default: Any = None) -> Any: return self.client.read(path, default)
[docs] def get(self, path: str, default: Any = None) -> Any: return self.client.get(path, default)
[docs] def read_many(self, *paths: str) -> dict[str, Any]: return self.client.read_many(*paths)
[docs] def get_values(self, *paths: str) -> dict[str, Any]: return self.client.read_many(*paths)
[docs] def keys(self) -> list[str]: return self.client.keys()
[docs] def values(self) -> dict[str, Any]: return self.client.values()
[docs] def flat_values(self) -> dict[str, Any]: return self.client.flat_values()
[docs] def meta(self) -> dict[str, Any]: return self.client.meta()
[docs] def snapshot(self) -> dict[str, Any]: return self.client.snapshot()
[docs] def get_snapshot(self) -> dict[str, Any]: return self.client.snapshot()
[docs] def messages(self, category: str | None = None, limit: int | None = None) -> list[MessageRecord]: return self.client.messages(category=category, limit=limit)
[docs] def latest_message(self, category: str | None = None) -> MessageRecord | None: return self.client.latest_message(category=category)
[docs] def latest_error_code(self) -> str | None: latest = self.latest_error or {} value = latest.get("error_code") or latest.get("code") if value: return str(value) detected = self.client.read("messages.latest_detected_error.detected_error_code") if detected: return str(detected) raw_detected = self.client.read("diagnostics.latest_raw_detected_error.suspected_error_code") if raw_detected: return str(raw_detected) return None
[docs] def lookup_error(self, code: str | None = None) -> dict[str, Any] | None: target = code or self.latest_error_code() if not target: return None return self.client.lookup_error(target)
[docs] def latest_error_details(self) -> dict[str, Any] | None: return self.client.latest_error_details()
[docs] def latest_message_error_details(self) -> dict[str, Any] | None: return self.client.latest_message_error_details()
[docs] def latest_raw_error_details(self) -> dict[str, Any] | None: return self.client.latest_raw_error_details()
[docs] def get_variable(self, name: str, default: Any = None) -> Any: return self.client.get_variable(name, default)
[docs] def subscribe(self, callback: Any) -> None: self.client.subscribe(callback)
[docs] def unsubscribe(self, callback: Any) -> None: self.client.unsubscribe(callback)
[docs] def watch(self, *paths: str, interval: float = 0.5) -> Iterator[dict[str, Any]]: return self.client.watch(*paths, interval=interval)
[docs] def watch_values(self, *paths: str, interval: float = 0.5, iterations: int | None = None) -> Iterator[dict[str, Any]]: count = 0 while iterations is None or count < max(0, int(iterations)): if paths: yield self.client.read_many(*paths) else: yield self.client.snapshot() count += 1 if iterations is not None and count >= max(0, int(iterations)): break self.client.wait_for_update(timeout=max(0.0, interval))
[docs] def get_robot_state(self) -> dict[str, Any]: return { "connected": self.connected, "status": self.status, "packet_count": self.packet_count, "robot_mode": self.robot_mode, "safety_mode": self.safety_mode, "tcp_pose": self.tcp_pose, "joint_positions": self.joint_positions, "global_variables": self.global_variables, "latest_error": self.latest_error, }
URClientAPI = ClientInterfaceAPI ClientInterface = ClientInterfaceAPI URClientMonitor = URClient __all__ = [ "AppSettings", "KeyMessage", "MessageRecord", "URClient", "URClientAPI", "ClientInterfaceAPI", "ClientInterface", "URClientMonitor", "interface_options", "interface_label", "interface_note", "supports_key_messages", "resolve_port", "pretty_json", "format_value", "format_robot_error_code", "detect_ur_error_code", "normalize_ur_error_code", "ascii_preview", "extract_printable_strings", "raw_packet_summary", ]