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 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 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 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",
]