#!/usr/bin/python3
# @file
# @ingroup uatools_tools
# @copyright
#   SPDX-FileCopyrightText: 2026 European Southern Observatory (ESO)
#   SPDX-License-Identifier: LGPL-3.0-only

"""
Minimal OPC UA test server using asyncua.

Tree (names compact for clarity):
- system (State, Init, Enable, Disable, Reset)
  - subsys1
    - fcs
      - lamp1, lamp2 (State, Init, Enable, Disable, Reset, On(intensity:int), Off)
      - motor1, motor2 (State, Init, Enable, Disable, Reset, MoveAbs(position:float, velocity:float), MoveVel(velocity:float), Stop)
      - sensor1, sensor2 (State, Init, Enable, Disable, Reset)
    - cameras
      - camera1, camera2 (State, Init, Enable, Disable, Reset, StartAcq, StopAcq)
  - subsys2 (same structure as subsys1)

Auth: optional, hardcoded user/user_pswd.
Size: -s/--size s|m|l|h controls namespace size:
  s/small: 5 subsystems, ~4K nodes
  m/medium: 10 subsystems, ~15K nodes
  l/large: 10 subsystems, ~100K nodes
  h/huge: 10 subsystems, ~200K nodes
CLI: UaTstServer -i/--ip <addr> -p/--port <port> [-a/--authentication] [-s/--size SIZE] [-l/--log-level LEVEL|logger:LEVEL] [--list-loggers]
"""

import argparse
import asyncio
import logging
import random
import re
import time
from collections import defaultdict
from contextlib import contextmanager
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from typing import Callable, Dict, List, Optional

from asyncua import Server, ua
try:
    import asyncua as _asyncua
    ASYNCUA_VERSION = getattr(_asyncua, "__version__", None)
except Exception:
    ASYNCUA_VERSION = None


LOGGER = logging.getLogger("uatools.pysrv")


# --- Startup profiler --------------------------------------------------------
# Opt-in via --profile-startup. No-op when disabled. Used to diagnose the
# (~4 minute) build time of `-s L`/`-s H` namespaces and identify which phase
# (server.init, address-space construction, _make_writable, etc.) dominates.
#
# The profiler tracks two things:
#   - Named phases (time spent in a `with PROFILER.phase('name'):` block).
#   - Per-builder counts/times via `PROFILER.tick(builder_name, secs)`.

class _Profiler:
    def __init__(self):
        self.enabled: bool = False
        self.phases: list = []        # (name, secs)
        self.builder_secs: dict = defaultdict(float)
        self.builder_count: dict = defaultdict(int)

    @contextmanager
    def phase(self, name: str):
        if not self.enabled:
            yield
            return
        t0 = time.perf_counter()
        try:
            yield
        finally:
            self.phases.append((name, time.perf_counter() - t0))

    def tick(self, name: str, secs: float) -> None:
        if not self.enabled:
            return
        self.builder_secs[name] += secs
        self.builder_count[name] += 1

    def report(self) -> None:
        if not self.enabled:
            return
        total_phase = sum(secs for _, secs in self.phases)
        print("\n=== UaTstServer startup profile =========================", flush=True)
        print(f"Total wall time across phases: {total_phase:8.3f}s", flush=True)
        print("Phases:", flush=True)
        for name, secs in self.phases:
            pct = (100.0 * secs / total_phase) if total_phase > 0 else 0.0
            print(f"  {name:38s} {secs:8.3f}s  ({pct:5.1f}%)", flush=True)
        if self.builder_count:
            total_builder = sum(self.builder_secs.values())
            print(f"\nPer-builder breakdown (total {total_builder:.3f}s, "
                  f"{sum(self.builder_count.values())} devices):", flush=True)
            rows = sorted(self.builder_secs.items(), key=lambda kv: kv[1], reverse=True)
            for name, secs in rows:
                n = self.builder_count[name]
                avg_ms = 1000.0 * secs / n if n else 0.0
                pct = (100.0 * secs / total_builder) if total_builder > 0 else 0.0
                print(f"  {name:24s} n={n:5d}  total={secs:7.3f}s  "
                      f"avg={avg_ms:6.2f}ms  ({pct:5.1f}%)", flush=True)
        print("==========================================================\n", flush=True)


PROFILER = _Profiler()


# --- Size constants ----------------------------------------------------------
SIZE_SMALL = "s"
SIZE_MEDIUM = "m"
SIZE_LARGE = "l"
SIZE_HUGE = "h"


def parse_size(size_str: str) -> str:
    """Parse size string (case-insensitive) and return normalized size constant."""
    size_lower = size_str.lower() if size_str else ""
    if size_lower in ("s", "small"):
        return SIZE_SMALL
    if size_lower in ("m", "medium"):
        return SIZE_MEDIUM
    if size_lower in ("l", "large"):
        return SIZE_LARGE
    if size_lower in ("h", "huge"):
        return SIZE_HUGE
    return SIZE_SMALL  # Default


# --- Auth --------------------------------------------------------------------
#
# Hard-coded credentials when -a/--authentication is on:
#   username: "user"
#   password: "user_pswd"
#
# asyncua 1.1.x changed the user_manager contract: where pre-1.1 it was an
# async callable `(username, password) -> bool`, in 1.1.x the iserver calls
# `user_manager.get_user(iserver, username, password, certificate)` and
# treats `None` as "reject". The old callable returned a coroutine (truthy),
# so every non-empty username was silently accepted. Implementing the new
# API class fixes the wrong-creds-rejected case.
#
# Cross-version note: asyncua 1.1.4 was the first release to export `User`
# and `UserRole` from `asyncua.crypto.permission_rules`. On older versions
# (e.g. 1.1.0 on platform25) those names don't exist, so we can't build a
# `User` object to return from get_user. We treat that as "auth not
# supported here": `-a` is refused with a clear error, anonymous launches
# keep working.
HARD_USER = "user"
HARD_PASS = "user_pswd"


# Probe the auth types once at import time. AUTH_SUPPORTED is consulted
# from main() so `-a` can be rejected before we even open a socket.
try:
    from asyncua.crypto.permission_rules import User as _AsyncuaUser
    from asyncua.crypto.permission_rules import UserRole as _AsyncuaUserRole
    AUTH_SUPPORTED = True
except ImportError:
    _AsyncuaUser = None
    _AsyncuaUserRole = None
    AUTH_SUPPORTED = False


class HardcodedUserManager:
    """UserManager implementing the asyncua 1.1.x get_user() contract.

    Returns an authenticated `User` for the hardcoded credential pair when
    auth is on, and `None` for everything else (which iserver translates
    into BadUserAccessDenied). Anonymous sessions never reach here —
    they're gated by the configured UserTokenPolicy on the endpoint.

    Only constructible on asyncua >= 1.1.4 (when `User`/`UserRole` were
    exported from `permission_rules`); see AUTH_SUPPORTED.
    """

    def __init__(self):
        if not AUTH_SUPPORTED:
            raise RuntimeError(
                "HardcodedUserManager requires asyncua >= 1.1.4 "
                "(User / UserRole not exported from "
                "asyncua.crypto.permission_rules)"
            )
        self._User = _AsyncuaUser
        self._UserRole = _AsyncuaUserRole

    def get_user(self, iserver, username=None, password=None, certificate=None):
        ok = (username == HARD_USER and password == HARD_PASS)
        LOGGER.info("Auth attempt user=%r password=%s -> %s",
                    username, "<set>" if password else "<empty>",
                    "OK" if ok else "REJECTED")
        if ok:
            return self._User(role=self._UserRole.User)
        return None


def build_user_token_policies(auth_enabled: bool) -> List[ua.UserTokenPolicy]:
    """Create explicit UserTokenPolicy entries so endpoints advertise supported identity tokens."""
    none_uri = "http://opcfoundation.org/UA/SecurityPolicy#None"
    policies: List[ua.UserTokenPolicy] = []

    if not auth_enabled:
        anon = ua.UserTokenPolicy()
        anon.TokenType = ua.UserTokenType.Anonymous
        anon.PolicyId = "anonymous"
        anon.SecurityPolicyUri = none_uri
        policies.append(anon)

    username = ua.UserTokenPolicy()
    username.TokenType = ua.UserTokenType.UserName
    username.PolicyId = "username"
    username.SecurityPolicyUri = none_uri
    policies.append(username)

    username_basic = ua.UserTokenPolicy()
    username_basic.TokenType = ua.UserTokenType.UserName
    username_basic.PolicyId = "username_basic256"
    username_basic.SecurityPolicyUri = none_uri
    policies.append(username_basic)
    return policies


# --- Device state helpers ----------------------------------------------------
STATE_NOTOP_NOTREADY = "NotOperational/NotReady"
STATE_NOTOP_READY = "NotOperational/Ready"
STATE_OP_IDLE = "Operational/Idle"


SUBSTATE_DEFAULTS = {
    "lamp": {
        STATE_NOTOP_NOTREADY: "Off",
        STATE_NOTOP_READY: "Off",
        STATE_OP_IDLE: "Off",
    },
    "motor": {
        STATE_NOTOP_NOTREADY: "Standstill",
        STATE_NOTOP_READY: "Standstill",
        STATE_OP_IDLE: "Standstill",
    },
    "sensor": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Reading",
        STATE_OP_IDLE: "Reading",
    },
    "camera": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    "subsystem": {
        STATE_NOTOP_NOTREADY: "",
        STATE_NOTOP_READY: "",
        STATE_OP_IDLE: "",
    },
    "system": {
        STATE_NOTOP_NOTREADY: "",
        STATE_NOTOP_READY: "",
        STATE_OP_IDLE: "",
    },
    # Motion/Positioning devices
    "shutter": {
        STATE_NOTOP_NOTREADY: "Closed",
        STATE_NOTOP_READY: "Closed",
        STATE_OP_IDLE: "Closed",
    },
    "filter_wheel": {
        STATE_NOTOP_NOTREADY: "Unknown",
        STATE_NOTOP_READY: "Unknown",
        STATE_OP_IDLE: "Unknown",
    },
    "focus": {
        STATE_NOTOP_NOTREADY: "Standstill",
        STATE_NOTOP_READY: "Standstill",
        STATE_OP_IDLE: "Standstill",
    },
    "derotator": {
        STATE_NOTOP_NOTREADY: "Stopped",
        STATE_NOTOP_READY: "Stopped",
        STATE_OP_IDLE: "Stopped",
    },
    "piezo": {
        STATE_NOTOP_NOTREADY: "Standstill",
        STATE_NOTOP_READY: "Standstill",
        STATE_OP_IDLE: "Standstill",
    },
    "hexapod": {
        STATE_NOTOP_NOTREADY: "Parked",
        STATE_NOTOP_READY: "Parked",
        STATE_OP_IDLE: "Parked",
    },
    # Optical devices
    "adc": {
        STATE_NOTOP_NOTREADY: "Parked",
        STATE_NOTOP_READY: "Parked",
        STATE_OP_IDLE: "Parked",
    },
    "polarizer": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    "grating": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    "mirror": {
        STATE_NOTOP_NOTREADY: "Flat",
        STATE_NOTOP_READY: "Flat",
        STATE_OP_IDLE: "Flat",
    },
    # Thermal/Environment devices
    "cooler": {
        STATE_NOTOP_NOTREADY: "Off",
        STATE_NOTOP_READY: "Off",
        STATE_OP_IDLE: "Off",
    },
    "heater": {
        STATE_NOTOP_NOTREADY: "Off",
        STATE_NOTOP_READY: "Off",
        STATE_OP_IDLE: "Off",
    },
    "vacuum": {
        STATE_NOTOP_NOTREADY: "Vented",
        STATE_NOTOP_READY: "Vented",
        STATE_OP_IDLE: "Vented",
    },
    "chiller": {
        STATE_NOTOP_NOTREADY: "Off",
        STATE_NOTOP_READY: "Off",
        STATE_OP_IDLE: "Off",
    },
    # Detectors
    "detector": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    "spectrograph": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    "wavefront_sensor": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    # Power/Safety devices
    "power_supply": {
        STATE_NOTOP_NOTREADY: "Disabled",
        STATE_NOTOP_READY: "Disabled",
        STATE_OP_IDLE: "Disabled",
    },
    "interlock": {
        STATE_NOTOP_NOTREADY: "Disarmed",
        STATE_NOTOP_READY: "Disarmed",
        STATE_OP_IDLE: "Disarmed",
    },
    "plc": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Idle",
        STATE_OP_IDLE: "Idle",
    },
    # Telemetry devices
    "weather_station": {
        STATE_NOTOP_NOTREADY: "Idle",
        STATE_NOTOP_READY: "Monitoring",
        STATE_OP_IDLE: "Monitoring",
    },
    "gps": {
        STATE_NOTOP_NOTREADY: "NoFix",
        STATE_NOTOP_READY: "NoFix",
        STATE_OP_IDLE: "Fixed",
    },
    "logger": {
        STATE_NOTOP_NOTREADY: "Stopped",
        STATE_NOTOP_READY: "Stopped",
        STATE_OP_IDLE: "Stopped",
    },
}


@dataclass
class DeviceContext:
    name: str
    device_type: str = "device"
    parent: Optional["DeviceContext"] = None
    children: List["DeviceContext"] = field(default_factory=list)
    obj: Optional[ua.NodeId] = None
    state_var: Optional[ua.NodeId] = None
    enabled_var: Optional[ua.NodeId] = None
    status_var: Optional[ua.NodeId] = None
    extra_vars: Dict[str, ua.NodeId] = field(default_factory=dict)
    target_position: Optional[float] = None
    velocity: float = 0.0
    telemetry: Dict[str, float] = field(default_factory=dict)
    noise_rng: random.Random = field(default_factory=random.Random)
    path: str = ""

    async def set_state(self, state: str):
        if self.state_var:
            await self.state_var.write_value(state)
        if self.status_var:
            await self.status_var.write_value(state)

    async def set_enabled(self, enabled: bool):
        if self.enabled_var:
            await self.enabled_var.write_value(enabled)


# --- State helpers -----------------------------------------------------------
def compose_state(base_state: str, substate: Optional[str]) -> str:
    if substate:
        return f"{base_state}/{substate}"
    return base_state


def base_state_of(state_val: str) -> str:
    parts = str(state_val).split("/")
    return "/".join(parts[:2]) if len(parts) >= 2 else str(state_val)


def substate_of(state_val: str) -> str:
    parts = str(state_val).split("/")
    return parts[2] if len(parts) >= 3 else ""


def default_substate(ctx: DeviceContext, base_state: str) -> str:
    return SUBSTATE_DEFAULTS.get(ctx.device_type, {}).get(base_state, "")


def derived_base_from_children(child_bases):
    if not child_bases:
        return STATE_NOTOP_NOTREADY
    normalized = [base_state_of(b) for b in child_bases]
    if STATE_NOTOP_NOTREADY in normalized:
        return STATE_NOTOP_NOTREADY
    if STATE_NOTOP_READY in normalized:
        return STATE_NOTOP_READY
    return STATE_OP_IDLE


async def update_derived_state(ctx: DeviceContext):
    """Derive parent's base state from immediate children."""
    if not ctx or not ctx.children:
        return
    child_bases = []
    for child in ctx.children:
        if not child.state_var:
            continue
        child_state = await child.state_var.read_value()
        child_bases.append(base_state_of(child_state))
    if not child_bases:
        return
    derived = derived_base_from_children(child_bases)
    await ctx.set_state(compose_state(derived, default_substate(ctx, derived)))
    if ctx.enabled_var:
        await ctx.set_enabled(derived == STATE_OP_IDLE)


async def propagate_up(ctx: Optional[DeviceContext]):
    """After a leaf state change, refresh ancestor derived states."""
    current = ctx
    while current:
        await update_derived_state(current)
        current = current.parent


_TYPED_BAD_INVALID_STATE = None
try:
    # asyncua >= 0.9 lives here
    from asyncua.ua.uaerrors._auto import BadInvalidState as _TYPED_BAD_INVALID_STATE
except ImportError:
    try:
        # older layout
        from asyncua.ua.uaerrors import BadInvalidState as _TYPED_BAD_INVALID_STATE
    except ImportError:
        _TYPED_BAD_INVALID_STATE = None


def _raise_bad_invalid_state():
    """Raise BadInvalidState so the client sees the **actual** status
    code, not asyncua's catch-all BadUnexpectedError.

    Why both raise paths matter: asyncua's address_space dispatcher
    inspects the *type* of the raised exception. The typed
    `BadInvalidState` subclass (when available) is recognised and its
    status code is preserved end-to-end. The generic
    `UaStatusCodeError(ua.StatusCodes.BadInvalidState)` falls through
    a broader handler in some versions and gets normalised to
    `BadUnexpectedError(0x80010000)`. We prefer the typed exception
    when the install provides it."""
    if _TYPED_BAD_INVALID_STATE is not None:
        raise _TYPED_BAD_INVALID_STATE()
    raise ua.UaStatusCodeError(ua.StatusCodes.BadInvalidState)


async def ensure_base_state(ctx: DeviceContext, expected_base: str):
    """Raise BadInvalidState if the current base state doesn't match."""
    if not ctx or not ctx.state_var:
        return
    current = await ctx.state_var.read_value()
    if base_state_of(current) != expected_base:
        # OPC UA precondition failure: the state-machine transition
        # requires a specific base state. WARN here (not in
        # transition_scope) because we now actually KNOW the
        # transition is illegal — we have both expected and actual
        # state — so the message is informative, not noise. The
        # successful path stays at INFO via transition_scope's own
        # log. method_wrapper catches the raise, logs the bookend
        # INFO, and re-raises so the status code reaches the client.
        LOGGER.warning(
            "Transition rejected on %s: state is %s, required base %s",
            getattr(ctx, "path", ctx.name),
            current,
            expected_base,
        )
        _raise_bad_invalid_state()


async def transition_scope(ctx: DeviceContext, target_base: str, allowed_from: Optional[str], enable: Optional[bool]):
    """Apply a base state transition to ctx and all descendants, enforcing preconditions."""
    if not ctx:
        return

    # INFO, not WARNING: this line fires on every transition attempt,
    # including the successful ones, so WARN here would cry wolf. The
    # failure path is already announced separately by method_wrapper's
    # "Returning error response ..." INFO line, which carries the
    # status-code class (BadInvalidState).
    LOGGER.info(
        "Transition request: %s -> %s (allowed_from=%s)",
        getattr(ctx, "path", ctx.name),
        target_base,
        allowed_from or "any",
    )

    async def _apply(node: DeviceContext):
        if allowed_from:
            await ensure_base_state(node, allowed_from)
        resolved_sub = default_substate(node, target_base)
        await node.set_state(compose_state(target_base, resolved_sub))
        if enable is not None:
            await node.set_enabled(enable)
        for child in node.children:
            await _apply(child)

    await _apply(ctx)
    await propagate_up(ctx.parent)


# --- Node creation -----------------------------------------------------------

# AccessLevel / UserAccessLevel bit mask:
#   bit 0 (1) = CurrentRead
#   bit 1 (2) = CurrentWrite
#   bit 2 (4) = HistoryRead
#   bit 3 (8) = HistoryWrite
# asyncua's `node.set_writable()` toggles the write bit on AccessLevel but
# does NOT touch UserAccessLevel — and OPC UA clients (including
# UaExplorer's strict-mode write button) prefer UserAccessLevel when
# present. To avoid clients showing "Writable: Not Available" or
# disabling the Write field, set BOTH attributes explicitly here.

_ACCESS_READ = 0x01
_ACCESS_RW = 0x03


async def _set_access_level(var, level: int) -> None:
    """Set both AccessLevel and UserAccessLevel on a Variable node so
    OPC UA clients reliably see whether the node is readable / writable.
    asyncua-side `set_writable()` alone misses UserAccessLevel."""
    await var.write_attribute(
        ua.AttributeIds.AccessLevel,
        ua.DataValue(ua.Variant(level, ua.VariantType.Byte)),
    )
    await var.write_attribute(
        ua.AttributeIds.UserAccessLevel,
        ua.DataValue(ua.Variant(level, ua.VariantType.Byte)),
    )


async def _make_writable(var) -> None:
    """Mark a Variable readable + writable for both AccessLevel and
    UserAccessLevel. Replaces ``set_writable()`` everywhere in this
    module — the asyncua call only sets AccessLevel and leaves
    UserAccessLevel unset, which trips up strict clients."""
    await var.set_writable()
    await _set_access_level(var, _ACCESS_RW)


async def _make_readonly(var) -> None:
    """Mark a Variable explicitly readable, not writable. Useful for
    sensor-style read-only nodes so the AccessLevel attribute is
    deterministic (otherwise newer asyncua releases sometimes leave
    UserAccessLevel unset, causing UaExplorer's strict-mode Write
    button to be disabled even though the value is readable)."""
    await _set_access_level(var, _ACCESS_READ)


async def add_state_vars(obj, name: str, node_id_prefix: str, ns_idx: int, device_type: str) -> DeviceContext:
    ctx = DeviceContext(name=name, device_type=device_type, obj=obj, path=node_id_prefix)
    initial_sub = SUBSTATE_DEFAULTS.get(device_type, {}).get(STATE_NOTOP_NOTREADY, "")
    initial_state = compose_state(STATE_NOTOP_NOTREADY, initial_sub)
    ctx.state_var = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.State", ns_idx), "State", initial_state)
    ctx.enabled_var = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Enabled", ns_idx), "Enabled", False)
    ctx.status_var = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Status", ns_idx), "Status", initial_state)
    await _make_writable(ctx.state_var)
    await _make_writable(ctx.enabled_var)
    await _make_writable(ctx.status_var)

    # Dummy attributes for client-side testing of types + read/write modes.
    # Present on every device so client scripts have a known stable surface
    # regardless of device_type.
    #   Note (str, writable)     — text round-trip
    #   Counter (int, writable)  — int round-trip
    #   Scale (float, writable)  — float round-trip
    #   ReadOnlyCounter (int, READ-ONLY) — explicit non-writable, useful
    #     for verifying strict-mode write button disable + Writable=False
    ctx.extra_vars["note"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Note", ns_idx), "Note", "")
    ctx.extra_vars["counter"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Counter", ns_idx), "Counter", 0)
    ctx.extra_vars["scale"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Scale", ns_idx), "Scale", 1.0)
    ctx.extra_vars["readonly_counter"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.ReadOnlyCounter", ns_idx),
        "ReadOnlyCounter", 0)
    await _make_writable(ctx.extra_vars["note"])
    await _make_writable(ctx.extra_vars["counter"])
    await _make_writable(ctx.extra_vars["scale"])
    await _make_readonly(ctx.extra_vars["readonly_counter"])
    return ctx


# --- Method handlers ---------------------------------------------------------
def _status_code_from_exc(exc: ua.UaStatusCodeError) -> ua.StatusCode:
    """Extract a `ua.StatusCode` from a UaStatusCodeError-or-subclass.

    asyncua's typed exceptions (e.g. BadInvalidState) carry the code
    as a class attribute or via `args[0]`; the generic
    UaStatusCodeError takes it as the first ctor argument. Cover all
    three so we get the right wire-level code regardless of how the
    handler chose to raise."""
    code = getattr(exc, "code", None)
    if code is None and exc.args:
        code = exc.args[0]
    if code is None:
        return ua.StatusCode(ua.StatusCodes.BadUnexpectedError)
    if isinstance(code, ua.StatusCode):
        return code
    return ua.StatusCode(code)


def method_wrapper(fn: Callable):
    async def _wrapped(parent_ctx: DeviceContext, *args):
        ctx_path = getattr(parent_ctx, "path", getattr(parent_ctx, "name", "unknown"))
        # INFO bookends: a method invocation is an operator-visible
        # event worth surfacing in default-level logs. DEBUG carries
        # the (potentially noisy) argument list separately so
        # production logs stay readable while devs can opt-in.
        LOGGER.info("Received method invocation request: %s.%s()", ctx_path, fn.__name__)
        LOGGER.debug("Method call %s on %s args=%s", fn.__name__, ctx_path, args)
        try:
            result = await fn(parent_ctx, *args)
        except ua.UaStatusCodeError as exc:
            # asyncua 1.x's `_call` dispatcher in
            # asyncua.server.address_space wraps EVERY exception
            # (including typed UaStatusCodeError subclasses) as
            # BadUnexpectedError before sending the response —
            # see its `except Exception:` block. The only way to
            # preserve the actual status code on the wire is to
            # *return* a `ua.StatusCode` from the handler; asyncua
            # then assigns it verbatim to res.StatusCode (see
            # `elif isinstance(result, ua.StatusCode):` in _call).
            sc = _status_code_from_exc(exc)
            LOGGER.info(
                "Returning error response to method invocation request: %s.%s() -> %s",
                ctx_path, fn.__name__, exc.__class__.__name__)
            LOGGER.debug(
                "Method %s on %s rejected: %s (status 0x%08X)",
                fn.__name__, ctx_path, exc.__class__.__name__, int(sc.value))
            return sc
        except Exception as exc:  # noqa: BLE001
            LOGGER.exception("Method %s failed", fn.__name__)
            LOGGER.info(
                "Returning error response to method invocation request: %s.%s() -> BadInternalError",
                ctx_path, fn.__name__)
            return ua.StatusCode(ua.StatusCodes.BadInternalError)
        LOGGER.debug("Method %s on %s returned ok", fn.__name__, ctx_path)
        return result

    return _wrapped


def _as_float(value) -> float:
    """Coerce OPC UA Variant or plain value to float, raising BadTypeMismatch on failure."""
    try:
        if isinstance(value, ua.Variant):
            value = value.Value
        return float(value)
    except Exception as exc:  # noqa: BLE001
        raise ua.UaStatusCodeError(ua.StatusCodes.BadTypeMismatch) from exc


def bind_method(ctx: DeviceContext, fn: Callable):
    # No @method_wrapper here: `fn` is already wrapped at module scope
    # (every m_* coroutine has @method_wrapper). Wrapping again produced
    # duplicate "Method _bound on unknown rejected:" log lines on every
    # rejection — once from the inner wrap, once from the outer one.
    async def _bound(_ignored_parent, *args):
        return await fn(ctx, *args)

    return _bound


@method_wrapper
async def m_state(ctx: DeviceContext) -> str:
    val = await ctx.state_var.read_value()
    return (ua.Variant(val, ua.VariantType.String),)


@method_wrapper
async def m_init(ctx: DeviceContext) -> int:
    await transition_scope(ctx, STATE_NOTOP_READY, allowed_from=STATE_NOTOP_NOTREADY, enable=False)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_enable(ctx: DeviceContext) -> int:
    await transition_scope(ctx, STATE_OP_IDLE, allowed_from=STATE_NOTOP_READY, enable=True)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_disable(ctx: DeviceContext) -> int:
    await transition_scope(ctx, STATE_NOTOP_READY, allowed_from=STATE_OP_IDLE, enable=False)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_reset(ctx: DeviceContext) -> int:
    await transition_scope(ctx, STATE_NOTOP_NOTREADY, allowed_from=None, enable=False)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_on(ctx: DeviceContext, intensity: int) -> int:
    ctx.extra_vars["intensity"] and await ctx.extra_vars["intensity"].write_value(intensity)
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "On"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_off(ctx: DeviceContext) -> int:
    ctx.extra_vars.get("intensity") and await ctx.extra_vars["intensity"].write_value(0)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Off"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_move_abs(ctx: DeviceContext, position: float, velocity: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    position_val = _as_float(position)
    velocity_val = _as_float(velocity)
    current_pos = 0.0
    if "position" in ctx.extra_vars:
        current_pos = await ctx.extra_vars["position"].read_value()
        current_pos = _as_float(current_pos)

    # Move toward the target regardless of the sign provided by the caller.
    direction = 0.0
    if position_val > current_pos:
        direction = 1.0
    elif position_val < current_pos:
        direction = -1.0

    speed = max(abs(velocity_val), 0.1)
    ctx.target_position = position_val
    ctx.velocity = speed * direction if direction else 0.0

    ctx.extra_vars.get("target") and await ctx.extra_vars["target"].write_value(position_val)
    ctx.extra_vars.get("velocity") and await ctx.extra_vars["velocity"].write_value(ctx.velocity)
    if direction:
        LOGGER.info(
            "MoveAbs %s: pos=%.3f -> target=%.3f vel=%.3f",
            ctx.path,
            current_pos,
            position_val,
            ctx.velocity,
        )
        await ctx.set_state(compose_state(STATE_OP_IDLE, "Moving"))
    else:
        await ctx.set_state(compose_state(STATE_OP_IDLE, "Standstill"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_move_vel(ctx: DeviceContext, velocity: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    ctx.target_position = None
    vel = _as_float(velocity)
    ctx.velocity = vel
    ctx.extra_vars.get("velocity") and await ctx.extra_vars["velocity"].write_value(vel)
    if abs(vel) < 1e-6:
        await ctx.set_state(compose_state(STATE_OP_IDLE, "Standstill"))
    else:
        LOGGER.info("MoveVel %s: vel=%.3f", ctx.path, vel)
        await ctx.set_state(compose_state(STATE_OP_IDLE, "Moving"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop(ctx: DeviceContext) -> int:
    ctx.velocity = 0.0
    ctx.target_position = None
    ctx.extra_vars.get("velocity") and await ctx.extra_vars["velocity"].write_value(0.0)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Standstill"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_start_acq(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    # Seed frame rate randomly for each acquisition start
    fr = ctx.noise_rng.uniform(5.0, 60.0)
    if "frame_rate" in ctx.extra_vars:
        await ctx.extra_vars["frame_rate"].write_value(fr)
    if "frame_count" in ctx.extra_vars:
        await ctx.extra_vars["frame_count"].write_value(0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Acquiring"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop_acq(ctx: DeviceContext) -> int:
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Idle"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


# --- Method handlers for new device types ----------------------------------
# Motion/Positioning devices
@method_wrapper
async def m_open(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "position" in ctx.extra_vars:
        await ctx.extra_vars["position"].write_value(1.0)
    if "open_count" in ctx.extra_vars:
        count = await ctx.extra_vars["open_count"].read_value()
        await ctx.extra_vars["open_count"].write_value(count + 1)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Open"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_close(ctx: DeviceContext) -> int:
    if "position" in ctx.extra_vars:
        await ctx.extra_vars["position"].write_value(0.0)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Closed"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_select_slot(ctx: DeviceContext, slot: int) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    slot_val = int(slot)
    if "slot" in ctx.extra_vars:
        await ctx.extra_vars["slot"].write_value(slot_val)
    if "num_slots" in ctx.extra_vars:
        num_slots = await ctx.extra_vars["num_slots"].read_value()
        if slot_val < 0 or slot_val >= num_slots:
            raise ua.UaStatusCodeError(ua.StatusCodes.BadOutOfRange)
    if "filter_name" in ctx.extra_vars:
        await ctx.extra_vars["filter_name"].write_value(f"Filter{slot_val}")
    await ctx.set_state(compose_state(STATE_OP_IDLE, f"Slot{slot_val}"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_home(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "position" in ctx.extra_vars:
        await ctx.extra_vars["position"].write_value(0.0)
    if "slot" in ctx.extra_vars:
        await ctx.extra_vars["slot"].write_value(0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Homed"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_move_rel(ctx: DeviceContext, offset: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    offset_val = _as_float(offset)
    current_pos = 0.0
    if "position" in ctx.extra_vars:
        current_pos = await ctx.extra_vars["position"].read_value()
        current_pos = _as_float(current_pos)
    new_pos = current_pos + offset_val
    if "position" in ctx.extra_vars:
        await ctx.extra_vars["position"].write_value(new_pos)
    if "in_position" in ctx.extra_vars:
        await ctx.extra_vars["in_position"].write_value(True)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Moving"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_angle(ctx: DeviceContext, angle: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    angle_val = _as_float(angle)
    if "angle" in ctx.extra_vars:
        await ctx.extra_vars["angle"].write_value(angle_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Positioned"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_start_tracking(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "tracking" in ctx.extra_vars:
        await ctx.extra_vars["tracking"].write_value(True)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Tracking"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop_tracking(ctx: DeviceContext) -> int:
    if "tracking" in ctx.extra_vars:
        await ctx.extra_vars["tracking"].write_value(False)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Stopped"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_voltage(ctx: DeviceContext, voltage: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    voltage_val = _as_float(voltage)
    if "voltage" in ctx.extra_vars:
        await ctx.extra_vars["voltage"].write_value(voltage_val)
    if "range" in ctx.extra_vars:
        range_val = await ctx.extra_vars["range"].read_value()
        if abs(voltage_val) > range_val:
            raise ua.UaStatusCodeError(ua.StatusCodes.BadOutOfRange)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Positioned"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_move_to(ctx: DeviceContext, x: float, y: float, z: float, rx: float, ry: float, rz: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "x" in ctx.extra_vars:
        await ctx.extra_vars["x"].write_value(_as_float(x))
    if "y" in ctx.extra_vars:
        await ctx.extra_vars["y"].write_value(_as_float(y))
    if "z" in ctx.extra_vars:
        await ctx.extra_vars["z"].write_value(_as_float(z))
    if "rx" in ctx.extra_vars:
        await ctx.extra_vars["rx"].write_value(_as_float(rx))
    if "ry" in ctx.extra_vars:
        await ctx.extra_vars["ry"].write_value(_as_float(ry))
    if "rz" in ctx.extra_vars:
        await ctx.extra_vars["rz"].write_value(_as_float(rz))
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Moving"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_move_rel_6dof(ctx: DeviceContext, dx: float, dy: float, dz: float, drx: float, dry: float, drz: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    for axis, delta in [("x", dx), ("y", dy), ("z", dz), ("rx", drx), ("ry", dry), ("rz", drz)]:
        if axis in ctx.extra_vars:
            current = await ctx.extra_vars[axis].read_value()
            await ctx.extra_vars[axis].write_value(_as_float(current) + _as_float(delta))
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Moving"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_park(ctx: DeviceContext) -> int:
    # Park to default position (0,0,0,0,0,0 for hexapod, or 0 for others)
    if "x" in ctx.extra_vars:
        for axis in ["x", "y", "z", "rx", "ry", "rz"]:
            if axis in ctx.extra_vars:
                await ctx.extra_vars[axis].write_value(0.0)
    elif "position" in ctx.extra_vars:
        await ctx.extra_vars["position"].write_value(0.0)
    elif "angle" in ctx.extra_vars:
        await ctx.extra_vars["angle"].write_value(0.0)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Parked"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


# Optical devices
@method_wrapper
async def m_set_angles(ctx: DeviceContext, prism1_angle: float, prism2_angle: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "prism1_angle" in ctx.extra_vars:
        await ctx.extra_vars["prism1_angle"].write_value(_as_float(prism1_angle))
    if "prism2_angle" in ctx.extra_vars:
        await ctx.extra_vars["prism2_angle"].write_value(_as_float(prism2_angle))
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Positioned"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_auto_track(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "mode" in ctx.extra_vars:
        await ctx.extra_vars["mode"].write_value("AutoTrack")
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Tracking"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_mode(ctx: DeviceContext, mode: str) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    mode_str = str(mode)
    if "mode" in ctx.extra_vars:
        await ctx.extra_vars["mode"].write_value(mode_str)
    await ctx.set_state(compose_state(STATE_OP_IDLE, mode_str))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_select_order(ctx: DeviceContext, order: int) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    order_val = int(order)
    if "order" in ctx.extra_vars:
        await ctx.extra_vars["order"].write_value(order_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, f"Order{order_val}"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_wavelength(ctx: DeviceContext, wavelength: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    wavelength_val = _as_float(wavelength)
    if "wavelength" in ctx.extra_vars:
        await ctx.extra_vars["wavelength"].write_value(wavelength_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Configured"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_position_ttp(ctx: DeviceContext, tip: float, tilt: float, piston: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "tip" in ctx.extra_vars:
        await ctx.extra_vars["tip"].write_value(_as_float(tip))
    if "tilt" in ctx.extra_vars:
        await ctx.extra_vars["tilt"].write_value(_as_float(tilt))
    if "piston" in ctx.extra_vars:
        await ctx.extra_vars["piston"].write_value(_as_float(piston))
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Positioned"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_flatten(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "tip" in ctx.extra_vars:
        await ctx.extra_vars["tip"].write_value(0.0)
    if "tilt" in ctx.extra_vars:
        await ctx.extra_vars["tilt"].write_value(0.0)
    if "piston" in ctx.extra_vars:
        await ctx.extra_vars["piston"].write_value(0.0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Flat"))
    return (ua.Variant(0, ua.VariantType.Int16),)


# Thermal/Environment devices
@method_wrapper
async def m_set_temp(ctx: DeviceContext, setpoint: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    setpoint_val = _as_float(setpoint)
    if "setpoint" in ctx.extra_vars:
        await ctx.extra_vars["setpoint"].write_value(setpoint_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Regulating"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_cool_down(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "setpoint" in ctx.extra_vars:
        await ctx.extra_vars["setpoint"].write_value(-50.0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Cooling"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_warm_up(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "setpoint" in ctx.extra_vars:
        await ctx.extra_vars["setpoint"].write_value(20.0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Warming"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_turn_on(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "On"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_turn_off(ctx: DeviceContext) -> int:
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Off"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_pump(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "pump_status" in ctx.extra_vars:
        await ctx.extra_vars["pump_status"].write_value("Running")
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Pumping"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_vent(ctx: DeviceContext) -> int:
    if "pump_status" in ctx.extra_vars:
        await ctx.extra_vars["pump_status"].write_value("Stopped")
    if "valve_state" in ctx.extra_vars:
        await ctx.extra_vars["valve_state"].write_value("Open")
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Vented"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_isolate(ctx: DeviceContext) -> int:
    if "valve_state" in ctx.extra_vars:
        await ctx.extra_vars["valve_state"].write_value("Closed")
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Isolated"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_flow(ctx: DeviceContext, flow: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    flow_val = _as_float(flow)
    if "flow" in ctx.extra_vars:
        await ctx.extra_vars["flow"].write_value(flow_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Flowing"))
    return (ua.Variant(0, ua.VariantType.Int16),)


# Detector devices
@method_wrapper
async def m_configure(ctx: DeviceContext, exposure_time: float, gain: float, binning_x: int, binning_y: int) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "exposure_time" in ctx.extra_vars:
        await ctx.extra_vars["exposure_time"].write_value(_as_float(exposure_time))
    if "gain" in ctx.extra_vars:
        await ctx.extra_vars["gain"].write_value(_as_float(gain))
    if "binning_x" in ctx.extra_vars:
        await ctx.extra_vars["binning_x"].write_value(int(binning_x))
    if "binning_y" in ctx.extra_vars:
        await ctx.extra_vars["binning_y"].write_value(int(binning_y))
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Configured"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_expose(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Exposing"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_abort(ctx: DeviceContext) -> int:
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Idle"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_acquire(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Acquiring"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_start_loop(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Looping"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop_loop(ctx: DeviceContext) -> int:
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Idle"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_calibrate(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Calibrating"))
    return (ua.Variant(0, ua.VariantType.Int16),)


# Power/Safety devices
@method_wrapper
async def m_set_current(ctx: DeviceContext, current: float) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    current_val = _as_float(current)
    if "current" in ctx.extra_vars:
        await ctx.extra_vars["current"].write_value(current_val)
    if "voltage" in ctx.extra_vars and "current" in ctx.extra_vars:
        voltage = await ctx.extra_vars["voltage"].read_value()
        if "power" in ctx.extra_vars:
            await ctx.extra_vars["power"].write_value(_as_float(voltage) * current_val)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Regulating"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_ps_enable(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "output_enabled" in ctx.extra_vars:
        await ctx.extra_vars["output_enabled"].write_value(True)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Enabled"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_ps_disable(ctx: DeviceContext) -> int:
    if "output_enabled" in ctx.extra_vars:
        await ctx.extra_vars["output_enabled"].write_value(False)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Disabled"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_arm(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "armed" in ctx.extra_vars:
        await ctx.extra_vars["armed"].write_value(True)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Armed"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_disarm(ctx: DeviceContext) -> int:
    if "armed" in ctx.extra_vars:
        await ctx.extra_vars["armed"].write_value(False)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Disarmed"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_acknowledge(ctx: DeviceContext) -> int:
    if "triggered" in ctx.extra_vars:
        await ctx.extra_vars["triggered"].write_value(False)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Acknowledged"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_output(ctx: DeviceContext, output: int, value: bool) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    output_idx = int(output)
    value_bool = bool(value)
    if "outputs" in ctx.extra_vars:
        outputs = await ctx.extra_vars["outputs"].read_value() or []
        if output_idx < len(outputs):
            outputs[output_idx] = value_bool
            await ctx.extra_vars["outputs"].write_value(outputs)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Active"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_read_input(ctx: DeviceContext, input_idx: int) -> bool:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    idx = int(input_idx)
    if "inputs" in ctx.extra_vars:
        inputs = await ctx.extra_vars["inputs"].read_value() or []
        if idx < len(inputs):
            return (ua.Variant(inputs[idx], ua.VariantType.Boolean),)
    return (ua.Variant(False, ua.VariantType.Boolean),)


@method_wrapper
async def m_clear_alarm(ctx: DeviceContext, alarm_id: int) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    alarm_idx = int(alarm_id)
    if "alarms" in ctx.extra_vars:
        alarms = await ctx.extra_vars["alarms"].read_value() or []
        if alarm_idx < len(alarms):
            alarms[alarm_idx] = False
            await ctx.extra_vars["alarms"].write_value(alarms)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Cleared"))
    return (ua.Variant(0, ua.VariantType.Int16),)


# Telemetry devices
@method_wrapper
async def m_start_monitoring(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Monitoring"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop_monitoring(ctx: DeviceContext) -> int:
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Idle"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_sync(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Synced"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_get_position(ctx: DeviceContext) -> tuple:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    lat = await ctx.extra_vars["latitude"].read_value() if "latitude" in ctx.extra_vars else 0.0
    lon = await ctx.extra_vars["longitude"].read_value() if "longitude" in ctx.extra_vars else 0.0
    alt = await ctx.extra_vars["altitude"].read_value() if "altitude" in ctx.extra_vars else 0.0
    return (
        ua.Variant(lat, ua.VariantType.Double),
        ua.Variant(lon, ua.VariantType.Double),
        ua.Variant(alt, ua.VariantType.Double),
    )


@method_wrapper
async def m_start_recording(ctx: DeviceContext) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    if "recording" in ctx.extra_vars:
        await ctx.extra_vars["recording"].write_value(True)
    if "record_count" in ctx.extra_vars:
        await ctx.extra_vars["record_count"].write_value(0)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Recording"))
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_stop_recording(ctx: DeviceContext) -> int:
    if "recording" in ctx.extra_vars:
        await ctx.extra_vars["recording"].write_value(False)
    enabled = await ctx.enabled_var.read_value()
    target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
    await ctx.set_state(compose_state(target_base, "Stopped"))
    await propagate_up(ctx.parent)
    return (ua.Variant(0, ua.VariantType.Int16),)


@method_wrapper
async def m_set_path(ctx: DeviceContext, file_path: str) -> int:
    await ensure_base_state(ctx, STATE_OP_IDLE)
    path_str = str(file_path)
    if "file_path" in ctx.extra_vars:
        await ctx.extra_vars["file_path"].write_value(path_str)
    await ctx.set_state(compose_state(STATE_OP_IDLE, "Configured"))
    return (ua.Variant(0, ua.VariantType.Int16),)


# --- Helpers for method arguments -------------------------------------------
def make_arg(name: str, variant_type: ua.VariantType, description: str = "") -> ua.Argument:
    """Create a ua.Argument with name, data type, and optional description."""
    # Map VariantType enum to corresponding OPC UA DataType NodeId
    type_to_nodeid = {
        ua.VariantType.Boolean: ua.NodeId(ua.ObjectIds.Boolean),
        ua.VariantType.SByte: ua.NodeId(ua.ObjectIds.SByte),
        ua.VariantType.Byte: ua.NodeId(ua.ObjectIds.Byte),
        ua.VariantType.Int16: ua.NodeId(ua.ObjectIds.Int16),
        ua.VariantType.UInt16: ua.NodeId(ua.ObjectIds.UInt16),
        ua.VariantType.Int32: ua.NodeId(ua.ObjectIds.Int32),
        ua.VariantType.UInt32: ua.NodeId(ua.ObjectIds.UInt32),
        ua.VariantType.Int64: ua.NodeId(ua.ObjectIds.Int64),
        ua.VariantType.UInt64: ua.NodeId(ua.ObjectIds.UInt64),
        ua.VariantType.Float: ua.NodeId(ua.ObjectIds.Float),
        ua.VariantType.Double: ua.NodeId(ua.ObjectIds.Double),
        ua.VariantType.String: ua.NodeId(ua.ObjectIds.String),
        ua.VariantType.DateTime: ua.NodeId(ua.ObjectIds.DateTime),
        ua.VariantType.Guid: ua.NodeId(ua.ObjectIds.Guid),
        ua.VariantType.ByteString: ua.NodeId(ua.ObjectIds.ByteString),
    }
    arg = ua.Argument()
    arg.Name = name
    arg.DataType = type_to_nodeid.get(variant_type, ua.NodeId(ua.ObjectIds.BaseDataType))
    arg.ValueRank = -1  # Scalar
    arg.ArrayDimensions = []
    if description:
        arg.Description = ua.LocalizedText(description)
    return arg


# --- Builders ---------------------------------------------------------------
async def add_common_methods(obj, ctx: DeviceContext, ns_idx: int, node_id_prefix: str):
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StateMethod", ns_idx), "State", bind_method(ctx, m_state),
        [], [make_arg("state", ua.VariantType.String, "Current state")]
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Init", ns_idx), "Init", bind_method(ctx, m_init),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")]
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Enable", ns_idx), "Enable", bind_method(ctx, m_enable),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")]
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Disable", ns_idx), "Disable", bind_method(ctx, m_disable),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")]
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Reset", ns_idx), "Reset", bind_method(ctx, m_reset),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")]
    )


async def build_lamp(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="lamp")
    ctx.extra_vars["intensity"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Intensity", ns_idx), "Intensity", 0)
    await _make_writable(ctx.extra_vars["intensity"])
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.On", ns_idx),
        "On",
        bind_method(ctx, m_on),
        [make_arg("intensity", ua.VariantType.Int32, "Lamp intensity level")],
        [make_arg("status", ua.VariantType.Int16, "Return status")],
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Off", ns_idx), "Off", bind_method(ctx, m_off),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_motor(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="motor")
    ctx.extra_vars["position"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Position", ns_idx), "Position", 0.0)
    ctx.extra_vars["target"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Target", ns_idx), "Target", 0.0)
    ctx.extra_vars["velocity"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Velocity", ns_idx), "Velocity", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveAbs", ns_idx),
        "MoveAbs",
        bind_method(ctx, m_move_abs),
        [make_arg("position", ua.VariantType.Double, "Target position"),
         make_arg("velocity", ua.VariantType.Double, "Movement velocity")],
        [make_arg("status", ua.VariantType.Int16, "Return status")],
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveVel", ns_idx),
        "MoveVel",
        bind_method(ctx, m_move_vel),
        [make_arg("velocity", ua.VariantType.Double, "Target velocity")],
        [make_arg("status", ua.VariantType.Int16, "Return status")],
    )
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Stop", ns_idx), "Stop", bind_method(ctx, m_stop),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_sensor(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="sensor")

    # Deterministic-ish per-sensor noise so different sensors get distinct ranges.
    seed = hash(node_id_prefix) & 0xFFFFFFFF
    ctx.noise_rng = random.Random(seed)

    # Base values and allowed drift per tick
    ctx.telemetry = {
        "temperature": ctx.noise_rng.uniform(18.0, 28.0),
        "flow": ctx.noise_rng.uniform(5.0, 30.0),
        "pressure": ctx.noise_rng.uniform(1.0, 4.0),
    }

    ctx.extra_vars["reading"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Reading", ns_idx), "Reading", 0.0)
    ctx.extra_vars["temperature"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Temperature", ns_idx), "Temperature", ctx.telemetry["temperature"]
    )
    ctx.extra_vars["flow"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Flow", ns_idx), "Flow", ctx.telemetry["flow"]
    )
    ctx.extra_vars["pressure"] = await obj.add_variable(
        ua.NodeId(f"{node_id_prefix}.Pressure", ns_idx), "Pressure", ctx.telemetry["pressure"]
    )

    for var in ("reading", "temperature", "flow", "pressure"):
        await _make_writable(ctx.extra_vars[var])

    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    return ctx


async def build_camera(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="camera")
    ctx.extra_vars["frame_rate"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.FrameRate", ns_idx), "FrameRate", 0.0)
    ctx.extra_vars["frame_count"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.FrameCount", ns_idx), "FrameCount", 0)
    await _make_writable(ctx.extra_vars["frame_rate"])
    await _make_writable(ctx.extra_vars["frame_count"])
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StartAcq", ns_idx), "StartAcq", bind_method(ctx, m_start_acq),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StopAcq", ns_idx), "StopAcq", bind_method(ctx, m_stop_acq),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# --- New device builders ----------------------------------------------------
# Motion/Positioning devices
async def build_shutter(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="shutter")
    ctx.extra_vars["position"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Position", ns_idx), "Position", 0.0)
    ctx.extra_vars["open_count"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.OpenCount", ns_idx), "OpenCount", 0)
    await _make_writable(ctx.extra_vars["position"])
    await _make_writable(ctx.extra_vars["open_count"])
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Open", ns_idx), "Open", bind_method(ctx, m_open),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Close", ns_idx), "Close", bind_method(ctx, m_close),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_filter_wheel(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="filter_wheel")
    ctx.extra_vars["slot"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Slot", ns_idx), "Slot", 0)
    ctx.extra_vars["filter_name"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.FilterName", ns_idx), "FilterName", "Filter0")
    ctx.extra_vars["num_slots"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.NumSlots", ns_idx), "NumSlots", 8)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SelectSlot", ns_idx), "SelectSlot", bind_method(ctx, m_select_slot),
        [make_arg("slot", ua.VariantType.Int32, "Slot number")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Home", ns_idx), "Home", bind_method(ctx, m_home),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_focus(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="focus")
    ctx.extra_vars["position"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Position", ns_idx), "Position", 0.0)
    ctx.extra_vars["in_position"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.InPosition", ns_idx), "InPosition", True)
    ctx.extra_vars["backlash"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Backlash", ns_idx), "Backlash", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveAbs", ns_idx), "MoveAbs", bind_method(ctx, m_move_abs),
        [make_arg("position", ua.VariantType.Double, "Target position"),
         make_arg("velocity", ua.VariantType.Double, "Movement velocity")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveRel", ns_idx), "MoveRel", bind_method(ctx, m_move_rel),
        [make_arg("offset", ua.VariantType.Double, "Relative offset")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Home", ns_idx), "Home", bind_method(ctx, m_home),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_derotator(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="derotator")
    ctx.extra_vars["angle"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Angle", ns_idx), "Angle", 0.0)
    ctx.extra_vars["tracking"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Tracking", ns_idx), "Tracking", False)
    ctx.extra_vars["offset"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Offset", ns_idx), "Offset", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetAngle", ns_idx), "SetAngle", bind_method(ctx, m_set_angle),
        [make_arg("angle", ua.VariantType.Double, "Target angle")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StartTracking", ns_idx), "StartTracking", bind_method(ctx, m_start_tracking),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StopTracking", ns_idx), "StopTracking", bind_method(ctx, m_stop_tracking),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_piezo(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="piezo")
    ctx.extra_vars["position"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Position", ns_idx), "Position", 0.0)
    ctx.extra_vars["voltage"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Voltage", ns_idx), "Voltage", 0.0)
    ctx.extra_vars["range"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Range", ns_idx), "Range", 100.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveAbs", ns_idx), "MoveAbs", bind_method(ctx, m_move_abs),
        [make_arg("position", ua.VariantType.Double, "Target position"),
         make_arg("velocity", ua.VariantType.Double, "Movement velocity")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetVoltage", ns_idx), "SetVoltage", bind_method(ctx, m_set_voltage),
        [make_arg("voltage", ua.VariantType.Double, "Voltage")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_hexapod(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="hexapod")
    ctx.extra_vars["x"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.X", ns_idx), "X", 0.0)
    ctx.extra_vars["y"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Y", ns_idx), "Y", 0.0)
    ctx.extra_vars["z"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Z", ns_idx), "Z", 0.0)
    ctx.extra_vars["rx"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RX", ns_idx), "RX", 0.0)
    ctx.extra_vars["ry"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RY", ns_idx), "RY", 0.0)
    ctx.extra_vars["rz"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RZ", ns_idx), "RZ", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveTo", ns_idx), "MoveTo", bind_method(ctx, m_move_to),
        [make_arg("x", ua.VariantType.Double, "X position"),
         make_arg("y", ua.VariantType.Double, "Y position"),
         make_arg("z", ua.VariantType.Double, "Z position"),
         make_arg("rx", ua.VariantType.Double, "RX rotation"),
         make_arg("ry", ua.VariantType.Double, "RY rotation"),
         make_arg("rz", ua.VariantType.Double, "RZ rotation")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.MoveRel", ns_idx), "MoveRel", bind_method(ctx, m_move_rel_6dof),
        [make_arg("dx", ua.VariantType.Double, "X offset"),
         make_arg("dy", ua.VariantType.Double, "Y offset"),
         make_arg("dz", ua.VariantType.Double, "Z offset"),
         make_arg("drx", ua.VariantType.Double, "RX offset"),
         make_arg("dry", ua.VariantType.Double, "RY offset"),
         make_arg("drz", ua.VariantType.Double, "RZ offset")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Park", ns_idx), "Park", bind_method(ctx, m_park),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# Optical devices
async def build_adc(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="adc")
    ctx.extra_vars["prism1_angle"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Prism1Angle", ns_idx), "Prism1Angle", 0.0)
    ctx.extra_vars["prism2_angle"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Prism2Angle", ns_idx), "Prism2Angle", 0.0)
    ctx.extra_vars["mode"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Mode", ns_idx), "Mode", "Manual")
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetAngles", ns_idx), "SetAngles", bind_method(ctx, m_set_angles),
        [make_arg("prism1_angle", ua.VariantType.Double, "Prism 1 angle"),
         make_arg("prism2_angle", ua.VariantType.Double, "Prism 2 angle")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.AutoTrack", ns_idx), "AutoTrack", bind_method(ctx, m_auto_track),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Park", ns_idx), "Park", bind_method(ctx, m_park),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_polarizer(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="polarizer")
    ctx.extra_vars["angle"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Angle", ns_idx), "Angle", 0.0)
    ctx.extra_vars["mode"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Mode", ns_idx), "Mode", "Linear")
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetAngle", ns_idx), "SetAngle", bind_method(ctx, m_set_angle),
        [make_arg("angle", ua.VariantType.Double, "Polarizer angle")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetMode", ns_idx), "SetMode", bind_method(ctx, m_set_mode),
        [make_arg("mode", ua.VariantType.String, "Mode")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_grating(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="grating")
    ctx.extra_vars["angle"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Angle", ns_idx), "Angle", 0.0)
    ctx.extra_vars["order"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Order", ns_idx), "Order", 1)
    ctx.extra_vars["wavelength"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Wavelength", ns_idx), "Wavelength", 500.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SelectOrder", ns_idx), "SelectOrder", bind_method(ctx, m_select_order),
        [make_arg("order", ua.VariantType.Int32, "Diffraction order")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetWavelength", ns_idx), "SetWavelength", bind_method(ctx, m_set_wavelength),
        [make_arg("wavelength", ua.VariantType.Double, "Wavelength")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_mirror(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="mirror")
    ctx.extra_vars["tip"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Tip", ns_idx), "Tip", 0.0)
    ctx.extra_vars["tilt"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Tilt", ns_idx), "Tilt", 0.0)
    ctx.extra_vars["piston"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Piston", ns_idx), "Piston", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetPosition", ns_idx), "SetPosition", bind_method(ctx, m_set_position_ttp),
        [make_arg("tip", ua.VariantType.Double, "Tip angle"),
         make_arg("tilt", ua.VariantType.Double, "Tilt angle"),
         make_arg("piston", ua.VariantType.Double, "Piston position")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Flatten", ns_idx), "Flatten", bind_method(ctx, m_flatten),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# Thermal/Environment devices
async def build_cooler(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="cooler")
    ctx.extra_vars["setpoint"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Setpoint", ns_idx), "Setpoint", 20.0)
    ctx.extra_vars["actual_temp"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.ActualTemp", ns_idx), "ActualTemp", 20.0)
    ctx.extra_vars["power"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Power", ns_idx), "Power", 0.0)
    ctx.extra_vars["ramp_rate"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RampRate", ns_idx), "RampRate", 1.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetTemp", ns_idx), "SetTemp", bind_method(ctx, m_set_temp),
        [make_arg("setpoint", ua.VariantType.Double, "Temperature setpoint")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.CoolDown", ns_idx), "CoolDown", bind_method(ctx, m_cool_down),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.WarmUp", ns_idx), "WarmUp", bind_method(ctx, m_warm_up),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_heater(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="heater")
    ctx.extra_vars["setpoint"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Setpoint", ns_idx), "Setpoint", 20.0)
    ctx.extra_vars["actual_temp"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.ActualTemp", ns_idx), "ActualTemp", 20.0)
    ctx.extra_vars["power"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Power", ns_idx), "Power", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetTemp", ns_idx), "SetTemp", bind_method(ctx, m_set_temp),
        [make_arg("setpoint", ua.VariantType.Double, "Temperature setpoint")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.TurnOn", ns_idx), "TurnOn", bind_method(ctx, m_turn_on),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.TurnOff", ns_idx), "TurnOff", bind_method(ctx, m_turn_off),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_vacuum(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="vacuum")
    ctx.extra_vars["pressure"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Pressure", ns_idx), "Pressure", 1013.25)
    ctx.extra_vars["pump_status"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.PumpStatus", ns_idx), "PumpStatus", "Stopped")
    ctx.extra_vars["valve_state"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.ValveState", ns_idx), "ValveState", "Open")
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Pump", ns_idx), "Pump", bind_method(ctx, m_pump),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Vent", ns_idx), "Vent", bind_method(ctx, m_vent),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Isolate", ns_idx), "Isolate", bind_method(ctx, m_isolate),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_chiller(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="chiller")
    ctx.extra_vars["flow"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Flow", ns_idx), "Flow", 0.0)
    ctx.extra_vars["inlet_temp"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.InletTemp", ns_idx), "InletTemp", 20.0)
    ctx.extra_vars["outlet_temp"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.OutletTemp", ns_idx), "OutletTemp", 20.0)
    ctx.extra_vars["setpoint"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Setpoint", ns_idx), "Setpoint", 20.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetFlow", ns_idx), "SetFlow", bind_method(ctx, m_set_flow),
        [make_arg("flow", ua.VariantType.Double, "Flow rate")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetTemp", ns_idx), "SetTemp", bind_method(ctx, m_set_temp),
        [make_arg("setpoint", ua.VariantType.Double, "Temperature setpoint")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# Detector devices
async def build_detector(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="detector")
    ctx.extra_vars["exposure_time"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.ExposureTime", ns_idx), "ExposureTime", 1.0)
    ctx.extra_vars["gain"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Gain", ns_idx), "Gain", 1.0)
    ctx.extra_vars["binning_x"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.BinningX", ns_idx), "BinningX", 1)
    ctx.extra_vars["binning_y"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.BinningY", ns_idx), "BinningY", 1)
    ctx.extra_vars["roi"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.ROI", ns_idx), "ROI", "0,0,1024,1024")
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Configure", ns_idx), "Configure", bind_method(ctx, m_configure),
        [make_arg("exposure_time", ua.VariantType.Double, "Exposure time"),
         make_arg("gain", ua.VariantType.Double, "Gain"),
         make_arg("binning_x", ua.VariantType.Int32, "Binning X"),
         make_arg("binning_y", ua.VariantType.Int32, "Binning Y")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Expose", ns_idx), "Expose", bind_method(ctx, m_expose),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Abort", ns_idx), "Abort", bind_method(ctx, m_abort),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_spectrograph(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="spectrograph")
    ctx.extra_vars["wavelength_min"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.WavelengthMin", ns_idx), "WavelengthMin", 400.0)
    ctx.extra_vars["wavelength_max"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.WavelengthMax", ns_idx), "WavelengthMax", 800.0)
    ctx.extra_vars["resolution"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Resolution", ns_idx), "Resolution", 0.1)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Configure", ns_idx), "Configure", bind_method(ctx, m_configure),
        [make_arg("wavelength_min", ua.VariantType.Double, "Minimum wavelength"),
         make_arg("wavelength_max", ua.VariantType.Double, "Maximum wavelength"),
         make_arg("resolution", ua.VariantType.Double, "Resolution")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Acquire", ns_idx), "Acquire", bind_method(ctx, m_acquire),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_wavefront_sensor(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="wavefront_sensor")
    ctx.extra_vars["mode"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Mode", ns_idx), "Mode", "Static")
    ctx.extra_vars["frequency"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Frequency", ns_idx), "Frequency", 10.0)
    ctx.extra_vars["rms_error"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RMSError", ns_idx), "RMSError", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StartLoop", ns_idx), "StartLoop", bind_method(ctx, m_start_loop),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StopLoop", ns_idx), "StopLoop", bind_method(ctx, m_stop_loop),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Calibrate", ns_idx), "Calibrate", bind_method(ctx, m_calibrate),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# Power/Safety devices
async def build_power_supply(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="power_supply")
    ctx.extra_vars["voltage"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Voltage", ns_idx), "Voltage", 0.0)
    ctx.extra_vars["current"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Current", ns_idx), "Current", 0.0)
    ctx.extra_vars["power"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Power", ns_idx), "Power", 0.0)
    ctx.extra_vars["output_enabled"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.OutputEnabled", ns_idx), "OutputEnabled", False)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetVoltage", ns_idx), "SetVoltage", bind_method(ctx, m_set_voltage),
        [make_arg("voltage", ua.VariantType.Double, "Voltage")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetCurrent", ns_idx), "SetCurrent", bind_method(ctx, m_set_current),
        [make_arg("current", ua.VariantType.Double, "Current")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.EnableOutput", ns_idx), "EnableOutput", bind_method(ctx, m_ps_enable),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.DisableOutput", ns_idx), "DisableOutput", bind_method(ctx, m_ps_disable),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_interlock(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="interlock")
    ctx.extra_vars["armed"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Armed", ns_idx), "Armed", False)
    ctx.extra_vars["triggered"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Triggered", ns_idx), "Triggered", False)
    ctx.extra_vars["bypass"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Bypass", ns_idx), "Bypass", False)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Arm", ns_idx), "Arm", bind_method(ctx, m_arm),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Disarm", ns_idx), "Disarm", bind_method(ctx, m_disarm),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    # Note: Reset is already added by add_common_methods, so we don't add it again
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Acknowledge", ns_idx), "Acknowledge", bind_method(ctx, m_acknowledge),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_plc(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="plc")
    ctx.extra_vars["inputs"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Inputs", ns_idx), "Inputs", [False] * 8)
    ctx.extra_vars["outputs"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Outputs", ns_idx), "Outputs", [False] * 8)
    ctx.extra_vars["alarms"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Alarms", ns_idx), "Alarms", [False] * 8)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetOutput", ns_idx), "SetOutput", bind_method(ctx, m_set_output),
        [make_arg("output", ua.VariantType.Int32, "Output index"),
         make_arg("value", ua.VariantType.Boolean, "Output value")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.ReadInput", ns_idx), "ReadInput", bind_method(ctx, m_read_input),
        [make_arg("input_idx", ua.VariantType.Int32, "Input index")],
        [make_arg("value", ua.VariantType.Boolean, "Input value")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.ClearAlarm", ns_idx), "ClearAlarm", bind_method(ctx, m_clear_alarm),
        [make_arg("alarm_id", ua.VariantType.Int32, "Alarm ID")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


# Telemetry devices
async def build_weather_station(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="weather_station")
    seed = hash(node_id_prefix) & 0xFFFFFFFF
    ctx.noise_rng = random.Random(seed)
    ctx.telemetry = {
        "wind_speed": ctx.noise_rng.uniform(0.0, 20.0),
        "wind_dir": ctx.noise_rng.uniform(0.0, 360.0),
        "humidity": ctx.noise_rng.uniform(30.0, 80.0),
        "temperature": ctx.noise_rng.uniform(10.0, 30.0),
        "pressure": ctx.noise_rng.uniform(980.0, 1020.0),
        "rain": ctx.noise_rng.uniform(0.0, 5.0),
    }
    ctx.extra_vars["wind_speed"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.WindSpeed", ns_idx), "WindSpeed", ctx.telemetry["wind_speed"])
    ctx.extra_vars["wind_dir"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.WindDir", ns_idx), "WindDir", ctx.telemetry["wind_dir"])
    ctx.extra_vars["humidity"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Humidity", ns_idx), "Humidity", ctx.telemetry["humidity"])
    ctx.extra_vars["temperature"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Temperature", ns_idx), "Temperature", ctx.telemetry["temperature"])
    ctx.extra_vars["pressure"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Pressure", ns_idx), "Pressure", ctx.telemetry["pressure"])
    ctx.extra_vars["rain"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Rain", ns_idx), "Rain", ctx.telemetry["rain"])
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StartMonitoring", ns_idx), "StartMonitoring", bind_method(ctx, m_start_monitoring),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StopMonitoring", ns_idx), "StopMonitoring", bind_method(ctx, m_stop_monitoring),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


async def build_gps(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="gps")
    seed = hash(node_id_prefix) & 0xFFFFFFFF
    ctx.noise_rng = random.Random(seed)
    ctx.extra_vars["latitude"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Latitude", ns_idx), "Latitude", ctx.noise_rng.uniform(-90.0, 90.0))
    ctx.extra_vars["longitude"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Longitude", ns_idx), "Longitude", ctx.noise_rng.uniform(-180.0, 180.0))
    ctx.extra_vars["altitude"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Altitude", ns_idx), "Altitude", ctx.noise_rng.uniform(0.0, 5000.0))
    ctx.extra_vars["utc_time"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.UTCTime", ns_idx), "UTCTime", "")
    ctx.extra_vars["pps_offset"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.PPSOffset", ns_idx), "PPSOffset", 0.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.Sync", ns_idx), "Sync", bind_method(ctx, m_sync),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.GetPosition", ns_idx), "GetPosition", bind_method(ctx, m_get_position),
        [], [make_arg("latitude", ua.VariantType.Double, "Latitude"),
              make_arg("longitude", ua.VariantType.Double, "Longitude"),
              make_arg("altitude", ua.VariantType.Double, "Altitude")])
    return ctx


async def build_logger(ns_idx: int, parent, name: str, prefix: str) -> DeviceContext:
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="logger")
    ctx.extra_vars["recording"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Recording", ns_idx), "Recording", False)
    ctx.extra_vars["file_path"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.FilePath", ns_idx), "FilePath", "")
    ctx.extra_vars["record_count"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.RecordCount", ns_idx), "RecordCount", 0)
    ctx.extra_vars["rate"] = await obj.add_variable(ua.NodeId(f"{node_id_prefix}.Rate", ns_idx), "Rate", 1.0)
    for var in ctx.extra_vars.values():
        await _make_writable(var)
    await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StartRecording", ns_idx), "StartRecording", bind_method(ctx, m_start_recording),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.StopRecording", ns_idx), "StopRecording", bind_method(ctx, m_stop_recording),
        [], [make_arg("status", ua.VariantType.Int16, "Return status")])
    await obj.add_method(
        ua.NodeId(f"{node_id_prefix}.SetPath", ns_idx), "SetPath", bind_method(ctx, m_set_path),
        [make_arg("file_path", ua.VariantType.String, "File path")],
        [make_arg("status", ua.VariantType.Int16, "Return status")])
    return ctx


def _format_arg_name(arg) -> str:
    """Extract argument name from ua.Argument, handling bytes or string."""
    if hasattr(arg, 'Name') and arg.Name is not None:
        if isinstance(arg.Name, bytes):
            return arg.Name.decode('utf-8', errors='replace')
        return str(arg.Name)
    return '?'


def _format_arg_type(arg) -> str:
    """Extract a short type name from ua.Argument.DataType."""
    type_names = {
        1: 'Boolean', 2: 'SByte', 3: 'Byte', 4: 'Int16',
        5: 'UInt16', 6: 'Int32', 7: 'UInt32', 8: 'Int64',
        9: 'UInt64', 10: 'Float', 11: 'Double', 12: 'String',
        13: 'DateTime', 14: 'Guid', 15: 'ByteString'
    }
    try:
        type_id = arg.DataType.Identifier
        return type_names.get(type_id, f'Type{type_id}')
    except Exception:
        return '?'


async def list_namespace_tree(server: Server, ns_idx: int):
    """Print the namespace tree for the given namespace index and exit."""

    async def _get_method_signature(node) -> str:
        """Get method signature with input argument names and types."""
        input_args = []
        output_args = []
        try:
            children = await node.get_children()
            for child in children:
                bn = await child.read_browse_name()
                if bn.Name == "InputArguments":
                    input_args = await child.read_value() or []
                elif bn.Name == "OutputArguments":
                    output_args = await child.read_value() or []
        except Exception:
            pass

        # Format input args
        in_parts = []
        for arg in input_args:
            name = _format_arg_name(arg)
            typ = _format_arg_type(arg)
            in_parts.append(f"{name}:{typ}")

        # Format output args
        out_parts = []
        for arg in output_args:
            name = _format_arg_name(arg)
            typ = _format_arg_type(arg)
            out_parts.append(f"{name}:{typ}")

        in_str = ", ".join(in_parts) if in_parts else ""
        out_str = ", ".join(out_parts) if out_parts else "void"
        return f"({in_str}) -> {out_str}"

    async def _walk(node, depth: int):
        try:
            bn = await node.read_browse_name()
            name = bn.Name or str(bn)
        except Exception:
            name = str(node)
        try:
            nc = await node.read_node_class()
            nc_name = ua.NodeClass(nc).name
            if nc == ua.NodeClass.Method:
                sig = await _get_method_signature(node)
                name = f"{name}{sig} (Method)"
            else:
                name = f"{name} ({nc_name})"
        except Exception:
            pass
        print("  " * depth + name)
        try:
            children = await node.get_children()
        except Exception:
            children = []
        for child in children:
            nid = getattr(child, "nodeid", None)
            if nid and getattr(nid, "NamespaceIndex", None) != ns_idx:
                continue
            await _walk(child, depth + 1)

    root = server.get_node(ua.NodeId("system", ns_idx))
    print(f"Namespace index {ns_idx}:")
    await _walk(root, 0)


async def build_subsystem(ns_idx: int, parent, name: str, prefix: str, stateful=True):
    node_id_prefix = f"{prefix}.{name}"
    obj = await parent.add_object(ua.NodeId(node_id_prefix, ns_idx), name)
    ctx = None
    if stateful:
        ctx = await add_state_vars(obj, name, node_id_prefix, ns_idx, device_type="subsystem")
        await add_common_methods(obj, ctx, ns_idx, node_id_prefix)
    return obj, ctx


# --- Simulation loop ---------------------------------------------------------
async def simulation_loop(devices: Dict[str, DeviceContext], period: float = 1.0):
    pos_dir = 1.0
    while True:
        for dev in devices.values():
            # Sensor telemetry and readings: only advance when Ready or Operational/Idle
            if dev.device_type == "sensor":
                try:
                    state_val = await dev.state_var.read_value()
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                base_state = base_state_of(state_val)
                if base_state == STATE_NOTOP_READY or base_state.startswith("Operational"):
                    if "reading" in dev.extra_vars:
                        val = (await dev.extra_vars["reading"].read_value()) + 0.1
                        await dev.extra_vars["reading"].write_value(val)

                    rng = getattr(dev, "noise_rng", None) or random.Random()

                    def bump(var_name: str, min_val: float, max_val: float, step: float):
                        if var_name not in dev.extra_vars:
                            return
                        current = dev.telemetry.get(var_name, (min_val + max_val) / 2)
                        delta = rng.uniform(-step, step)
                        new_val = max(min_val, min(max_val, current + delta))
                        dev.telemetry[var_name] = new_val
                        return dev.extra_vars[var_name].write_value(new_val)

                    # Temperature drift within a modest band
                    if (task := bump("temperature", 10.0, 40.0, 0.2)):
                        await task
                    # Flow varies a bit more; keep non-negative
                    if (task := bump("flow", 0.0, 60.0, 0.8)):
                        await task
                    # Pressure moves slowly in a smaller band
                    if (task := bump("pressure", 0.5, 6.0, 0.1)):
                        await task
            if "position" in dev.extra_vars:
                pos = await dev.extra_vars["position"].read_value()
                vel = dev.velocity
                if vel != 0.0:
                    pos += vel * period
                    await dev.extra_vars["position"].write_value(pos)
                    if dev.target_position is not None:
                        distance = dev.target_position - pos
                        step = abs(vel) * period
                        # Stop when we reach/overshoot the target.
                        if abs(distance) <= step or (distance * vel) <= 0:
                            pos = dev.target_position
                            await dev.extra_vars["position"].write_value(pos)
                            enabled = await dev.enabled_var.read_value()
                            target_base = STATE_OP_IDLE if enabled else STATE_NOTOP_READY
                            await dev.set_state(compose_state(target_base, "Standstill"))
                            LOGGER.info(
                                "MoveAbs %s: reached target %.3f (pos=%.3f), stopping",
                                dev.path,
                                dev.target_position,
                                pos,
                            )
                            dev.velocity = 0.0
                            dev.target_position = None
            if "frame_rate" in dev.extra_vars:
                state_val = await dev.state_var.read_value()
                base_state = base_state_of(state_val)
                sub_state = substate_of(state_val)
                if base_state == STATE_OP_IDLE and sub_state == "Acquiring":
                    fr = await dev.extra_vars["frame_rate"].read_value()
                    # Small jitter around current rate
                    jitter = getattr(dev, "noise_rng", random.Random()).uniform(-1.0, 1.0)
                    new_fr = max(1.0, fr + jitter)
                    await dev.extra_vars["frame_rate"].write_value(new_fr)
                    if "frame_count" in dev.extra_vars:
                        count = await dev.extra_vars["frame_count"].read_value()
                        await dev.extra_vars["frame_count"].write_value(count + int(max(new_fr * period, 1)))

            # Thermal devices: temperature drift toward setpoint
            if dev.device_type in ("cooler", "heater"):
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE:
                    if "setpoint" in dev.extra_vars and "actual_temp" in dev.extra_vars:
                        setpoint = await dev.extra_vars["setpoint"].read_value()
                        actual = await dev.extra_vars["actual_temp"].read_value()
                        ramp_rate = await dev.extra_vars["ramp_rate"].read_value() if "ramp_rate" in dev.extra_vars else 1.0
                        diff = setpoint - actual
                        step = min(abs(diff), ramp_rate * period) * (1.0 if diff > 0 else -1.0)
                        new_temp = actual + step
                        await dev.extra_vars["actual_temp"].write_value(new_temp)
                        if "power" in dev.extra_vars:
                            power = abs(diff) * 0.1  # Simple power model
                            await dev.extra_vars["power"].write_value(power)

            # Vacuum: pressure changes based on pump status
            if dev.device_type == "vacuum":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE and "pressure" in dev.extra_vars and "pump_status" in dev.extra_vars:
                    pump_status = await dev.extra_vars["pump_status"].read_value()
                    pressure = await dev.extra_vars["pressure"].read_value()
                    if pump_status == "Running":
                        # Pumping: decrease pressure
                        new_pressure = max(0.001, pressure - 10.0 * period)
                    else:
                        # Vented: increase toward atmospheric
                        new_pressure = min(1013.25, pressure + 5.0 * period)
                    await dev.extra_vars["pressure"].write_value(new_pressure)

            # Chiller: flow and temperature updates
            if dev.device_type == "chiller":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE:
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    if "inlet_temp" in dev.extra_vars:
                        inlet = await dev.extra_vars["inlet_temp"].read_value()
                        delta = rng.uniform(-0.1, 0.1)
                        await dev.extra_vars["inlet_temp"].write_value(inlet + delta)
                    if "outlet_temp" in dev.extra_vars and "setpoint" in dev.extra_vars:
                        setpoint = await dev.extra_vars["setpoint"].read_value()
                        outlet = await dev.extra_vars["outlet_temp"].read_value()
                        diff = setpoint - outlet
                        step = min(abs(diff), 0.5 * period) * (1.0 if diff > 0 else -1.0)
                        await dev.extra_vars["outlet_temp"].write_value(outlet + step)

            # Shutter: position updates (already handled by position logic above)
            # Filter wheel: discrete position (no continuous update needed)
            # Focus/derotator/piezo: position updates (already handled by position logic above)
            # Hexapod: 6-DOF position updates (already handled by position logic for x, y, z)

            # Optical devices: angle updates when tracking
            if dev.device_type == "adc":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                    sub_state = substate_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                    sub_state = ""
                if base_state == STATE_OP_IDLE and sub_state == "Tracking":
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    if "prism1_angle" in dev.extra_vars:
                        angle1 = await dev.extra_vars["prism1_angle"].read_value()
                        await dev.extra_vars["prism1_angle"].write_value(angle1 + rng.uniform(-0.01, 0.01))
                    if "prism2_angle" in dev.extra_vars:
                        angle2 = await dev.extra_vars["prism2_angle"].read_value()
                        await dev.extra_vars["prism2_angle"].write_value(angle2 + rng.uniform(-0.01, 0.01))

            # Detector: exposure progress
            if dev.device_type == "detector":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                    sub_state = substate_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                    sub_state = ""
                if base_state == STATE_OP_IDLE and sub_state == "Exposing" and "exposure_time" in dev.extra_vars:
                    # Exposure progress would be tracked here if needed
                    pass

            # Spectrograph: wavelength updates
            if dev.device_type == "spectrograph":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE:
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    if "wavelength_min" in dev.extra_vars and "wavelength_max" in dev.extra_vars:
                        wl_min = await dev.extra_vars["wavelength_min"].read_value()
                        wl_max = await dev.extra_vars["wavelength_max"].read_value()
                        # Small drift
                        await dev.extra_vars["wavelength_min"].write_value(wl_min + rng.uniform(-0.01, 0.01))
                        await dev.extra_vars["wavelength_max"].write_value(wl_max + rng.uniform(-0.01, 0.01))

            # Wavefront sensor: RMS error updates when in loop
            if dev.device_type == "wavefront_sensor":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                    sub_state = substate_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                    sub_state = ""
                if base_state == STATE_OP_IDLE and sub_state == "Looping" and "rms_error" in dev.extra_vars:
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    rms = await dev.extra_vars["rms_error"].read_value()
                    new_rms = max(0.0, rms + rng.uniform(-0.01, 0.01))
                    await dev.extra_vars["rms_error"].write_value(new_rms)

            # Weather station: telemetry updates
            if dev.device_type == "weather_station":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                    sub_state = substate_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                    sub_state = ""
                if base_state == STATE_OP_IDLE and sub_state == "Monitoring":
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    def bump_weather(var_name: str, min_val: float, max_val: float, step: float):
                        if var_name not in dev.extra_vars:
                            return
                        current = dev.telemetry.get(var_name, (min_val + max_val) / 2)
                        delta = rng.uniform(-step, step)
                        new_val = max(min_val, min(max_val, current + delta))
                        dev.telemetry[var_name] = new_val
                        return dev.extra_vars[var_name].write_value(new_val)

                    if (task := bump_weather("wind_speed", 0.0, 50.0, 0.5)):
                        await task
                    if (task := bump_weather("wind_dir", 0.0, 360.0, 2.0)):
                        await task
                    if (task := bump_weather("humidity", 0.0, 100.0, 0.5)):
                        await task
                    if (task := bump_weather("temperature", -10.0, 40.0, 0.2)):
                        await task
                    if (task := bump_weather("pressure", 950.0, 1050.0, 0.5)):
                        await task
                    if (task := bump_weather("rain", 0.0, 10.0, 0.1)):
                        await task

            # GPS: position and time updates
            if dev.device_type == "gps":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE:
                    rng = getattr(dev, "noise_rng", None) or random.Random()
                    if "latitude" in dev.extra_vars:
                        lat = await dev.extra_vars["latitude"].read_value()
                        await dev.extra_vars["latitude"].write_value(lat + rng.uniform(-0.0001, 0.0001))
                    if "longitude" in dev.extra_vars:
                        lon = await dev.extra_vars["longitude"].read_value()
                        await dev.extra_vars["longitude"].write_value(lon + rng.uniform(-0.0001, 0.0001))
                    if "utc_time" in dev.extra_vars:
                        await dev.extra_vars["utc_time"].write_value(datetime.now(timezone.utc).isoformat())

            # Logger: record count increments when recording
            if dev.device_type == "logger":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                    sub_state = substate_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                    sub_state = ""
                if base_state == STATE_OP_IDLE and sub_state == "Recording":
                    if "record_count" in dev.extra_vars and "rate" in dev.extra_vars:
                        count = await dev.extra_vars["record_count"].read_value()
                        rate = await dev.extra_vars["rate"].read_value()
                        await dev.extra_vars["record_count"].write_value(count + int(max(rate * period, 1)))

            # Power supply: voltage/current updates when enabled
            if dev.device_type == "power_supply":
                try:
                    state_val = await dev.state_var.read_value()
                    base_state = base_state_of(state_val)
                except Exception:
                    state_val = STATE_NOTOP_NOTREADY
                    base_state = base_state_of(state_val)
                if base_state == STATE_OP_IDLE and "output_enabled" in dev.extra_vars:
                    enabled = await dev.extra_vars["output_enabled"].read_value()
                    if enabled:
                        rng = getattr(dev, "noise_rng", None) or random.Random()
                        if "voltage" in dev.extra_vars:
                            voltage = await dev.extra_vars["voltage"].read_value()
                            await dev.extra_vars["voltage"].write_value(voltage + rng.uniform(-0.01, 0.01))
                        if "current" in dev.extra_vars and "voltage" in dev.extra_vars:
                            voltage = await dev.extra_vars["voltage"].read_value()
                            current = await dev.extra_vars["current"].read_value()
                            if "power" in dev.extra_vars:
                                await dev.extra_vars["power"].write_value(voltage * current)
        await asyncio.sleep(period)


# --- Server setup ------------------------------------------------------------
def configure_logging(level_arg: str):
    fmt_kwargs = dict(
        level=logging.WARNING,
        format="%(asctime)s.%(msecs)03d %(levelname)s:%(name)s:%(message)s",
        datefmt="%Y-%m-%dT%H:%M:%S",
    )
    base_level = logging.WARNING
    min_override = logging.CRITICAL
    overrides = []
    if level_arg:
        for token in level_arg.split(","):
            token = token.strip()
            if not token:
                continue
            if ":" in token:
                logger_name, lvl = token.split(":", 1)
                lvl_val = getattr(logging, lvl.upper(), logging.INFO)
                overrides.append((logger_name, lvl_val))
                min_override = min(min_override, lvl_val)
            else:
                base_level = getattr(logging, token.upper(), logging.INFO)
    effective_root = min(base_level, min_override) if min_override != logging.CRITICAL else base_level
    fmt_kwargs["level"] = effective_root
    fmt_kwargs["force"] = True  # Override any existing configuration
    logging.basicConfig(**fmt_kwargs)
    # Set our own logger to the base level
    LOGGER.setLevel(base_level)
    # Default quiet for noisy asyncua components unless explicitly overridden.
    # NB: asyncua.server.server is bumped to ERROR to suppress the
    # "No encrypting policy available, password may get transferred in
    # plaintext" WARNING that fires every server init. UaTstServer is
    # localhost / dev / CI only, doesn't carry real credentials, and
    # never installs a certificate — the warning is correct but useless
    # in this context and floods test output. Users who pass
    # --log-level asyncua.server.server:WARNING re-enable it.
    default_quiet = {
        "asyncua": logging.WARNING,
        "asyncua.server.uaprocessor": logging.WARNING,
        "asyncua.server.server": logging.ERROR,
        # Note: we used to silence asyncua.server.address_space at
        # CRITICAL because its dispatcher logged a traceback on every
        # handler raise. method_wrapper now *returns* a ua.StatusCode
        # instead of raising (asyncua 1.x's `_call` wraps any raise as
        # BadUnexpectedError, ignoring the actual UaStatusCodeError
        # subclass — returning is the only way to preserve the code),
        # so the traceback path no longer fires and we keep
        # address_space at its default level so genuine asyncua errors
        # stay visible.
    }
    for lname, lvl in default_quiet.items():
        logging.getLogger(lname).setLevel(lvl)
    # Apply explicit overrides
    for lname, lvl in overrides:
        logging.getLogger(lname).setLevel(lvl)


def _parse_version_tuple(version_str: Optional[str]) -> Optional[tuple]:
    if not version_str:
        return None
    parts = []
    for token in str(version_str).split("."):
        try:
            parts.append(int(token))
        except Exception:
            break
    return tuple(parts) if parts else None


# --- Structured-Variable fixture --------------------------------------------
# Regression fixture for the "structured Variable not traversed" bug. On real
# TwinCAT/Siemens PLC servers, DataBlocks are exposed as *Variable* nodes that
# have CHILD nodes (their fields) via HasComponent — a "structured Variable".
# UaExplorer/UaShell/ualib originally recursed only into Object nodes, so every
# node under such a Variable was silently hidden. This fixture reproduces that
# shape so the traversal fix has local test coverage (the normal device tree is
# all Objects-with-leaf-Variables and cannot trigger the bug).
#
# Shape added under /Objects:
#   DataBlockVar               (Variable)   <- a "DataBlock" as a Variable
#     Field1                   (Variable)   <- child of a Variable
#     Field2                   (Variable)
#     SubStruct                (Variable)   <- nested structured Variable
#       DeepField              (Variable)   <- child of a child Variable
async def build_structured_variable_fixture(ns_idx: int, objects_node):
    """Add a Variable node that has child Variable nodes (a mini PLC DataBlock),
    so recursive-browse fixes that descend into Variables have a test target."""
    db = await objects_node.add_variable(
        ua.NodeId("DataBlockVar", ns_idx), "DataBlockVar", 0)
    await db.add_variable(ua.NodeId("DataBlockVar.Field1", ns_idx), "Field1", 1)
    await db.add_variable(ua.NodeId("DataBlockVar.Field2", ns_idx), "Field2", 2)
    sub = await db.add_variable(
        ua.NodeId("DataBlockVar.SubStruct", ns_idx), "SubStruct", 0)
    await sub.add_variable(
        ua.NodeId("DataBlockVar.SubStruct.DeepField", ns_idx), "DeepField", 3)
    LOGGER.info("Structured-Variable fixture added: DataBlockVar with child fields")


async def build_hierarchy(ns_idx: int, root, system_ctx: DeviceContext, size: str, devices: Dict[str, DeviceContext]):
    """Build device hierarchy based on size parameter. Scales by devices per subsystem, max 10 subsystems."""
    size_norm = parse_size(size)

    # Determine number of subsystems and device configuration based on size.
    #
    # Target node counts (approximate):
    #   small  – 5 subsys, ~4K nodes   (~42 devices/subsys, 6 containers)
    #   medium – 10 subsys, ~15K nodes  (~96 devices/subsys, 8 containers)
    #   large  – 10 subsys, ~100K nodes (~695 devices/subsys, 8 containers)
    #   huge   – 10 subsys, ~200K nodes (~1330 devices/subsys, 8 containers)

    if size_norm == SIZE_SMALL:
        num_subsystems = 5
        # Small: ~42 devices/subsys × 5 subsys → ~4K nodes
        subsystem_configs = [
            ("fcs",
             [f"lamp{i}" for i in range(1, 6)] + [f"motor{i}" for i in range(1, 6)] + [f"sensor{i}" for i in range(1, 6)],
             [build_lamp] * 5 + [build_motor] * 5 + [build_sensor] * 5),
            ("cameras",
             [f"camera{i}" for i in range(1, 6)],
             [build_camera] * 5),
            ("optics",
             [f"shutter{i}" for i in range(1, 3)] + [f"filter_wheel{i}" for i in range(1, 3)] + [f"focus{i}" for i in range(1, 3)] + ["derotator1"],
             [build_shutter] * 2 + [build_filter_wheel] * 2 + [build_focus] * 2 + [build_derotator]),
            ("environment",
             [f"cooler{i}" for i in range(1, 3)] + [f"heater{i}" for i in range(1, 3)] + ["vacuum1", "chiller1"],
             [build_cooler] * 2 + [build_heater] * 2 + [build_vacuum] + [build_chiller]),
            ("detectors",
             [f"detector{i}" for i in range(1, 3)] + [f"spectrograph{i}" for i in range(1, 3)] + ["wavefront_sensor1"],
             [build_detector] * 2 + [build_spectrograph] * 2 + [build_wavefront_sensor]),
            ("telemetry",
             ["weather_station1", "gps1", "logger1", "logger2"],
             [build_weather_station, build_gps, build_logger, build_logger]),
        ]

    elif size_norm == SIZE_MEDIUM:
        num_subsystems = 10
        # Medium: ~96 devices/subsys × 10 subsys → ~15K nodes
        subsystem_configs = [
            ("fcs",
             [f"lamp{i}" for i in range(1, 10)] + [f"motor{i}" for i in range(1, 10)] + [f"sensor{i}" for i in range(1, 10)],
             [build_lamp] * 9 + [build_motor] * 9 + [build_sensor] * 9),
            ("cameras",
             [f"camera{i}" for i in range(1, 10)],
             [build_camera] * 9),
            ("optics",
             [f"shutter{i}" for i in range(1, 4)] + [f"filter_wheel{i}" for i in range(1, 4)] + [f"focus{i}" for i in range(1, 4)] + [f"derotator{i}" for i in range(1, 3)] + [f"piezo{i}" for i in range(1, 3)] + ["hexapod1"],
             [build_shutter] * 3 + [build_filter_wheel] * 3 + [build_focus] * 3 + [build_derotator] * 2 + [build_piezo] * 2 + [build_hexapod]),
            ("environment",
             [f"cooler{i}" for i in range(1, 4)] + [f"heater{i}" for i in range(1, 4)] + [f"vacuum{i}" for i in range(1, 3)] + [f"chiller{i}" for i in range(1, 3)],
             [build_cooler] * 3 + [build_heater] * 3 + [build_vacuum] * 2 + [build_chiller] * 2),
            ("detectors",
             [f"detector{i}" for i in range(1, 4)] + [f"spectrograph{i}" for i in range(1, 4)] + [f"wavefront_sensor{i}" for i in range(1, 3)],
             [build_detector] * 3 + [build_spectrograph] * 3 + [build_wavefront_sensor] * 2),
            ("optics2",
             [f"adc{i}" for i in range(1, 4)] + [f"polarizer{i}" for i in range(1, 4)] + [f"grating{i}" for i in range(1, 4)] + [f"mirror{i}" for i in range(1, 4)],
             [build_adc] * 3 + [build_polarizer] * 3 + [build_grating] * 3 + [build_mirror] * 3),
            ("power",
             [f"power_supply{i}" for i in range(1, 4)] + [f"interlock{i}" for i in range(1, 3)] + [f"plc{i}" for i in range(1, 3)],
             [build_power_supply] * 3 + [build_interlock] * 2 + [build_plc] * 2),
            ("telemetry",
             [f"weather_station{i}" for i in range(1, 4)] + [f"gps{i}" for i in range(1, 4)] + [f"logger{i}" for i in range(1, 4)],
             [build_weather_station] * 3 + [build_gps] * 3 + [build_logger] * 3),
        ]

    elif size_norm == SIZE_LARGE:
        num_subsystems = 10
        # Large: ~695 devices/subsys × 10 subsys → ~100K nodes
        subsystem_configs = [
            ("fcs",
             [f"lamp{i}" for i in range(1, 51)] + [f"motor{i}" for i in range(1, 51)] + [f"sensor{i}" for i in range(1, 51)],
             [build_lamp] * 50 + [build_motor] * 50 + [build_sensor] * 50),
            ("cameras",
             [f"camera{i}" for i in range(1, 61)],
             [build_camera] * 60),
            ("optics",
             [f"shutter{i}" for i in range(1, 31)] + [f"filter_wheel{i}" for i in range(1, 31)] + [f"focus{i}" for i in range(1, 31)] + [f"derotator{i}" for i in range(1, 21)] + [f"piezo{i}" for i in range(1, 21)] + [f"hexapod{i}" for i in range(1, 11)],
             [build_shutter] * 30 + [build_filter_wheel] * 30 + [build_focus] * 30 + [build_derotator] * 20 + [build_piezo] * 20 + [build_hexapod] * 10),
            ("environment",
             [f"cooler{i}" for i in range(1, 26)] + [f"heater{i}" for i in range(1, 26)] + [f"vacuum{i}" for i in range(1, 21)] + [f"chiller{i}" for i in range(1, 21)],
             [build_cooler] * 25 + [build_heater] * 25 + [build_vacuum] * 20 + [build_chiller] * 20),
            ("detectors",
             [f"detector{i}" for i in range(1, 26)] + [f"spectrograph{i}" for i in range(1, 26)] + [f"wavefront_sensor{i}" for i in range(1, 21)],
             [build_detector] * 25 + [build_spectrograph] * 25 + [build_wavefront_sensor] * 20),
            ("optics2",
             [f"adc{i}" for i in range(1, 21)] + [f"polarizer{i}" for i in range(1, 21)] + [f"grating{i}" for i in range(1, 21)] + [f"mirror{i}" for i in range(1, 21)],
             [build_adc] * 20 + [build_polarizer] * 20 + [build_grating] * 20 + [build_mirror] * 20),
            ("power",
             [f"power_supply{i}" for i in range(1, 26)] + [f"interlock{i}" for i in range(1, 21)] + [f"plc{i}" for i in range(1, 21)],
             [build_power_supply] * 25 + [build_interlock] * 20 + [build_plc] * 20),
            ("telemetry",
             [f"weather_station{i}" for i in range(1, 11)] + [f"gps{i}" for i in range(1, 16)] + [f"logger{i}" for i in range(1, 16)],
             [build_weather_station] * 10 + [build_gps] * 15 + [build_logger] * 15),
        ]

    elif size_norm == SIZE_HUGE:
        num_subsystems = 10
        # Huge: ~1330 devices/subsys × 10 subsys → ~200K nodes
        subsystem_configs = [
            ("fcs",
             [f"lamp{i}" for i in range(1, 81)] + [f"motor{i}" for i in range(1, 81)] + [f"sensor{i}" for i in range(1, 81)],
             [build_lamp] * 80 + [build_motor] * 80 + [build_sensor] * 80),
            ("cameras",
             [f"camera{i}" for i in range(1, 101)],
             [build_camera] * 100),
            ("optics",
             [f"shutter{i}" for i in range(1, 51)] + [f"filter_wheel{i}" for i in range(1, 51)] + [f"focus{i}" for i in range(1, 51)] + [f"derotator{i}" for i in range(1, 31)] + [f"piezo{i}" for i in range(1, 31)] + [f"hexapod{i}" for i in range(1, 21)],
             [build_shutter] * 50 + [build_filter_wheel] * 50 + [build_focus] * 50 + [build_derotator] * 30 + [build_piezo] * 30 + [build_hexapod] * 20),
            ("environment",
             [f"cooler{i}" for i in range(1, 51)] + [f"heater{i}" for i in range(1, 51)] + [f"vacuum{i}" for i in range(1, 41)] + [f"chiller{i}" for i in range(1, 41)],
             [build_cooler] * 50 + [build_heater] * 50 + [build_vacuum] * 40 + [build_chiller] * 40),
            ("detectors",
             [f"detector{i}" for i in range(1, 51)] + [f"spectrograph{i}" for i in range(1, 51)] + [f"wavefront_sensor{i}" for i in range(1, 41)],
             [build_detector] * 50 + [build_spectrograph] * 50 + [build_wavefront_sensor] * 40),
            ("optics2",
             [f"adc{i}" for i in range(1, 41)] + [f"polarizer{i}" for i in range(1, 51)] + [f"grating{i}" for i in range(1, 51)] + [f"mirror{i}" for i in range(1, 51)],
             [build_adc] * 40 + [build_polarizer] * 50 + [build_grating] * 50 + [build_mirror] * 50),
            ("power",
             [f"power_supply{i}" for i in range(1, 51)] + [f"interlock{i}" for i in range(1, 41)] + [f"plc{i}" for i in range(1, 41)],
             [build_power_supply] * 50 + [build_interlock] * 40 + [build_plc] * 40),
            ("telemetry",
             [f"weather_station{i}" for i in range(1, 21)] + [f"gps{i}" for i in range(1, 51)] + [f"logger{i}" for i in range(1, 51)],
             [build_weather_station] * 20 + [build_gps] * 50 + [build_logger] * 50),
        ]

    else:
        # Default: same as small
        num_subsystems = 5
        subsystem_configs = [
            ("fcs",
             [f"lamp{i}" for i in range(1, 6)] + [f"motor{i}" for i in range(1, 6)] + [f"sensor{i}" for i in range(1, 6)],
             [build_lamp] * 5 + [build_motor] * 5 + [build_sensor] * 5),
            ("cameras",
             [f"camera{i}" for i in range(1, 6)],
             [build_camera] * 5),
            ("optics",
             [f"shutter{i}" for i in range(1, 3)] + [f"filter_wheel{i}" for i in range(1, 3)] + [f"focus{i}" for i in range(1, 3)] + ["derotator1"],
             [build_shutter] * 2 + [build_filter_wheel] * 2 + [build_focus] * 2 + [build_derotator]),
            ("environment",
             [f"cooler{i}" for i in range(1, 3)] + [f"heater{i}" for i in range(1, 3)] + ["vacuum1", "chiller1"],
             [build_cooler] * 2 + [build_heater] * 2 + [build_vacuum] + [build_chiller]),
            ("detectors",
             [f"detector{i}" for i in range(1, 3)] + [f"spectrograph{i}" for i in range(1, 3)] + ["wavefront_sensor1"],
             [build_detector] * 2 + [build_spectrograph] * 2 + [build_wavefront_sensor]),
            ("telemetry",
             ["weather_station1", "gps1", "logger1", "logger2"],
             [build_weather_station, build_gps, build_logger, build_logger]),
        ]

    # Build subsystems with containers
    for subsys_num in range(1, num_subsystems + 1):
        subsys_name = f"subsys{subsys_num:03d}"  # Zero-padded format: subsys001, subsys002, ...
        with PROFILER.phase(f"subsys[{subsys_name}]"):
            subsys_obj, subsys_ctx = await build_subsystem(ns_idx, root, subsys_name, "system", stateful=True)
            if subsys_ctx:
                subsys_ctx.parent = system_ctx
                system_ctx.children.append(subsys_ctx)

            for subsys_type, dev_names, builders in subsystem_configs:
                container, container_ctx = await build_subsystem(ns_idx, subsys_obj, subsys_type, f"system.{subsys_name}", stateful=True)
                if container_ctx:
                    container_ctx.parent = subsys_ctx
                    subsys_ctx.children.append(container_ctx)

                for dev_name, builder in zip(dev_names, builders):
                    dev_path = f"system.{subsys_name}.{subsys_type}.{dev_name}"
                    t0 = time.perf_counter() if PROFILER.enabled else 0.0
                    devices[dev_path] = await builder(ns_idx, container, dev_name, f"system.{subsys_name}.{subsys_type}")
                    if PROFILER.enabled:
                        PROFILER.tick(builder.__name__, time.perf_counter() - t0)
                    devices[dev_path].parent = container_ctx
                    container_ctx.children.append(devices[dev_path])


# ---------------------------------------------------------------------------
# XML profile mode
#
# When --xml-profile <file.xml> is given on the command line, the built-in
# namespace + simulation loop is REPLACED with whatever the XML defines.
# This lets the test server stand in for a real device whose NodeSet2 XML
# is already in hand (e.g. exported from UaExplorer's "Export Namespace").
#
# Without --xml-profile, the server behaves exactly as before.
#
# Optional --xml-profile-lib <python.module> binds user-written logic to
# the loaded namespace:
#   - Methods named "Foo" call `on_Foo(parent, *args) -> result_tuple`
#   - Writable Variables named "Bar" call `on_write_Bar(node, value)`
# Both callbacks are best-effort: missing implementations are quietly
# skipped, allowing partial logic libraries.
#
# --gen-profile-logic <out.py> generates a stub module from the XML and
# exits — the user fills in the bodies, builds it as an importable Python
# package, then passes it via --xml-profile-lib at run time.
# ---------------------------------------------------------------------------

async def _collect_xml_methods_and_writables(server) -> tuple:
    """After import_xml, walk the address space and pick out user-namespace
    methods + writable variables. Returns (methods, writables) where each
    is a list of dicts: {'name': str, 'browse_name': str, 'node': Node,
    'input_args': list of (name, type_label), 'parent': Node}."""
    methods: List[dict] = []
    writables: List[dict] = []

    async def _walk(node, depth=0):
        if depth > 32:
            return  # safety: deep cycles
        try:
            children = await node.get_children()
        except Exception:
            return
        for child in children:
            try:
                node_class = await child.read_node_class()
                browse = await child.read_browse_name()
                browse_name = browse.Name if hasattr(browse, "Name") else str(browse)
            except Exception:
                continue
            # Skip standard OPC UA namespace (ns=0) — only the user's
            # custom namespace is interesting for logic generation.
            try:
                nid = child.nodeid
                if getattr(nid, "NamespaceIndex", 0) == 0:
                    await _walk(child, depth + 1)
                    continue
            except Exception:
                pass

            if node_class == ua.NodeClass.Method:
                input_args = await _read_input_args(child)
                methods.append({
                    "name": browse_name,
                    "browse_name": browse_name,
                    "node": child,
                    "input_args": input_args,
                    "parent": node,
                })
            elif node_class == ua.NodeClass.Variable:
                try:
                    access = await child.read_attribute(ua.AttributeIds.UserAccessLevel)
                    val = access.Value.Value if hasattr(access, "Value") else None
                    if isinstance(val, int) and (val & 0x02):
                        writables.append({
                            "name": browse_name,
                            "browse_name": browse_name,
                            "node": child,
                            "parent": node,
                        })
                except Exception:
                    pass
            await _walk(child, depth + 1)

    await _walk(server.nodes.objects)
    return methods, writables


async def _read_input_args(method_node) -> List[tuple]:
    """Return [(name, type_label)] for a Method node, by reading its
    InputArguments property child. Empty list when none / unreadable."""
    try:
        for child in await method_node.get_children():
            browse = await child.read_browse_name()
            if (getattr(browse, "Name", "") or str(browse)) != "InputArguments":
                continue
            args_value = await child.read_value()
            if not args_value:
                return []
            out: List[tuple] = []
            for arg in args_value:
                name = getattr(arg, "Name", "") or getattr(arg, "name", "")
                dt = getattr(arg, "DataType", None)
                type_label = _xml_arg_type_label(dt)
                out.append((str(name), type_label))
            return out
    except Exception:
        pass
    return []


_XML_BUILTIN_TYPE_NAMES = {
    1: "bool", 2: "int", 3: "int", 4: "int", 5: "int",
    6: "int", 7: "int", 8: "int", 9: "int",
    10: "float", 11: "float",
    12: "str", 15: "bytes",
}


def _xml_arg_type_label(data_type) -> str:
    if data_type is None:
        return "Any"
    ident = getattr(data_type, "Identifier", None)
    ns_idx = getattr(data_type, "NamespaceIndex", None)
    if ns_idx == 0 and isinstance(ident, int):
        return _XML_BUILTIN_TYPE_NAMES.get(ident, "Any")
    return "Any"


def _safe_identifier(name: str) -> str:
    """Sanitise a BrowseName into a valid Python identifier suffix."""
    cleaned = re.sub(r"[^A-Za-z0-9_]", "_", str(name))
    if cleaned and cleaned[0].isdigit():
        cleaned = "_" + cleaned
    return cleaned or "Unnamed"


async def generate_profile_logic(xml_path: Path, out_path: Path) -> None:
    """Load ``xml_path`` into a throwaway Server, walk it to collect
    methods + writable variables, and emit a stub Python module to
    ``out_path``. Refuses to overwrite an existing file — the user
    must delete it manually to regenerate."""
    if out_path.exists():
        raise RuntimeError(
            f"{out_path} already exists; refusing to overwrite. "
            f"Delete the file and re-run --gen-profile-lib to regenerate."
        )

    server = Server()
    await server.init()
    server.set_endpoint("opc.tcp://127.0.0.1:14841")  # not bound; needed by init
    server.set_server_name("uatools-pysrv-genprofile")
    server.set_security_policy([ua.SecurityPolicyType.NoSecurity])
    await server.import_xml(str(xml_path))
    try:
        methods, writables = await _collect_xml_methods_and_writables(server)
    finally:
        # Don't `async with server:` — we never started serving.
        pass

    timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
    lines: List[str] = []
    lines.append(f"# Generated by UaTstServer --gen-profile-logic on {timestamp}")
    lines.append(f"# Source: {xml_path}")
    lines.append("#")
    lines.append("# Fill in the `on_<MethodName>` and `on_write_<VariableName>` bodies.")
    lines.append("# Functions left in their default form are no-ops (Methods return"
                 " None; writes are accepted unchanged).")
    lines.append("# Pass this module to UaTstServer via:")
    lines.append("#   UaTstServer --xml-profile <file.xml> --xml-profile-lib <module.path>")
    lines.append("")
    lines.append("from typing import Any")
    lines.append("")
    lines.append("")

    # Methods
    if methods:
        lines.append("# ----------------------------------------------------------")
        lines.append("# Method handlers — one per Method in the XML profile.")
        lines.append("# Each handler runs when a client calls the corresponding")
        lines.append("# OPC UA Method. Return value is sent back to the caller.")
        lines.append("# ----------------------------------------------------------")
        lines.append("")
        seen: set = set()
        for m in methods:
            fname = f"on_{_safe_identifier(m['name'])}"
            if fname in seen:
                # Same BrowseName in multiple parents — disambiguate
                # by appending a counter. Caller's bind step picks the
                # parent-specific name, but since the user owns the
                # logic file we can't auto-pick; emit the base name only
                # for the first occurrence and skip duplicates so the
                # file stays valid Python.
                continue
            seen.add(fname)
            params = ["parent"] + [
                f"{_safe_identifier(n)}: {t}" for n, t in m["input_args"]
            ]
            params_sig = ", ".join(params)
            arg_doc = ", ".join(
                f"{n}: {t}" for n, t in m["input_args"]
            ) or "<no args>"
            lines.append(f"async def {fname}({params_sig}) -> Any:")
            lines.append(f'    """Method {m["browse_name"]}({arg_doc}).')
            lines.append("")
            lines.append("    Fill in the body. Return value is sent back to the")
            lines.append("    OPC UA caller as the method's OutputArguments.")
            lines.append('    """')
            lines.append("    return None")
            lines.append("")
            lines.append("")

    # Writable variables
    if writables:
        lines.append("# ----------------------------------------------------------")
        lines.append("# Write hooks — one per writable Variable in the XML.")
        lines.append("# Called BEFORE the value is committed to the address space.")
        lines.append("# Return the (possibly modified) value to commit, or raise")
        lines.append("# to reject the write.")
        lines.append("# ----------------------------------------------------------")
        lines.append("")
        seen_w: set = set()
        for w in writables:
            fname = f"on_write_{_safe_identifier(w['name'])}"
            if fname in seen_w:
                continue
            seen_w.add(fname)
            lines.append(f"async def {fname}(node, value: Any) -> Any:")
            lines.append(f'    """Hook for writes to variable {w["browse_name"]}.')
            lines.append("")
            lines.append("    The default returns ``value`` unchanged, so writes")
            lines.append("    pass through to the address space as normal.")
            lines.append('    """')
            lines.append("    return value")
            lines.append("")
            lines.append("")

    if not methods and not writables:
        lines.append("# (No Methods or writable Variables found in the XML.)")
        lines.append("")

    out_path.parent.mkdir(parents=True, exist_ok=True)
    out_path.write_text("\n".join(lines), encoding="utf-8")
    LOGGER.info(
        "Wrote profile logic stub to %s (%d methods, %d writable variables)",
        out_path, len(methods), len(writables),
    )


def _import_profile_lib(module_path: str):
    """Import a dotted module path. Returns the module object, or raises."""
    import importlib
    return importlib.import_module(module_path)


class _UaTstSubHandler:
    """Subscription handler that dispatches DataChange events to the
    profile-lib's on_write_<VariableName> callbacks. Used for write
    hooks since asyncua's server-side Variable write callback API
    differs across versions; subscribing to our own server's data
    changes is the most-portable way."""

    def __init__(self, callbacks):
        self._callbacks = callbacks  # dict node_id -> async callable

    def datachange_notification(self, node, val, data):  # noqa: N802
        cb = self._callbacks.get(node.nodeid)
        if cb is None:
            return
        try:
            loop = asyncio.get_event_loop()
            loop.create_task(cb(node, val))
        except Exception as exc:
            LOGGER.warning("on_write callback failed for %s: %s", node, exc)


async def _bind_profile_logic(server, profile_lib) -> None:
    """Walk the imported namespace and attach callbacks where the
    user's profile_lib provides on_<Method> / on_write_<Var> functions.
    Missing callbacks are silently skipped — partial logic libs are OK."""
    methods, writables = await _collect_xml_methods_and_writables(server)
    bound_methods = 0
    bound_writes = 0

    # Methods — wrap user callback so it adapts to asyncua's expected
    # (parent, *args) signature and returns a list of Variants.
    for m in methods:
        fname = f"on_{_safe_identifier(m['name'])}"
        user_fn = getattr(profile_lib, fname, None)
        if user_fn is None:
            continue
        await _attach_method_callback(m["node"], user_fn)
        bound_methods += 1

    # Writes — collect into a single subscription handler.
    write_callbacks = {}
    for w in writables:
        fname = f"on_write_{_safe_identifier(w['name'])}"
        user_fn = getattr(profile_lib, fname, None)
        if user_fn is None:
            continue
        write_callbacks[w["node"].nodeid] = user_fn
        bound_writes += 1

    if write_callbacks:
        try:
            handler = _UaTstSubHandler(write_callbacks)
            sub = await server.create_subscription(100, handler)
            for nid in write_callbacks:
                await sub.subscribe_data_change(server.get_node(nid))
        except Exception as exc:
            LOGGER.warning("Failed to wire write hooks: %s", exc)

    LOGGER.info(
        "Profile logic bound: %d method handlers, %d write hooks",
        bound_methods, bound_writes,
    )


async def _attach_method_callback(method_node, user_fn) -> None:
    """Attach a user method handler to a Method node. asyncua's API is
    `link_method(node, callable)` where the callable receives
    (parent_nodeid, *args) and returns a list of Variants."""
    server = getattr(method_node, "server", None)

    async def _adapter(parent, *args):
        try:
            result = await user_fn(parent, *args)
        except Exception as exc:
            LOGGER.warning("Method %s raised: %s", method_node, exc)
            raise
        # Allow callbacks to return None or a single value; wrap into
        # the variant list asyncua's link_method expects.
        if result is None:
            return []
        if isinstance(result, (list, tuple)):
            return [r if isinstance(r, ua.Variant) else ua.Variant(r) for r in result]
        return [ua.Variant(result)]

    # asyncua exposes either `link_method` on the Node or has the user
    # call a server-level helper. Try both.
    if hasattr(method_node, "link_method"):
        method_node.link_method(_adapter)
        return
    if server is not None and hasattr(server, "link_method"):
        server.link_method(method_node, _adapter)
        return
    LOGGER.warning(
        "Could not attach handler for %s — no link_method API in this "
        "asyncua version", method_node,
    )


async def build_server(port: int, auth: bool, log_level: str, host: str, size: str = SIZE_SMALL):
    configure_logging(log_level)
    version_tuple = _parse_version_tuple(ASYNCUA_VERSION)
    if auth and version_tuple and version_tuple < (1, 1, 8):
        LOGGER.warning(
            "Authentication (-a) is not supported with asyncua %s (requires >=1.1.8); exiting.",
            ASYNCUA_VERSION,
        )
        raise SystemExit(2)
    server = Server()
    with PROFILER.phase("server.init"):
        await server.init()
    server.set_endpoint(f"opc.tcp://{host}:{port}")
    server.set_server_name("uatools-pysrv")
    server.set_security_policy([ua.SecurityPolicyType.NoSecurity])
    supports_identity_tokens = hasattr(server, "set_identity_tokens")
    # Configure user tokens explicitly:
    # - auth=True: username only, anonymous disabled.
    # - auth=False: allow anonymous (and also publish username for clients that try cached creds).
    if auth:
        token_ids_str = ["UserName"]
        token_ids_enum = [ua.UserTokenType.UserName]
    else:
        token_ids_str = ["Anonymous", "UserName"]
        token_ids_enum = [ua.UserTokenType.Anonymous, ua.UserTokenType.UserName]
    token_policies = build_user_token_policies(auth)

    def _try_set_security_ids(ids) -> bool:
        try:
            server.set_security_IDs(ids)
            return True
        except Exception as exc:
            LOGGER.debug("set_security_IDs failed for %s: %s", ids, exc)
            return False

    tokens_set = False
    if supports_identity_tokens:
        try:
            server.set_identity_tokens(token_policies)
            tokens_set = True
        except Exception as exc:
            LOGGER.debug("set_identity_tokens failed: %s", exc)
    if not tokens_set:
        tokens_set = _try_set_security_ids(token_ids_str) or _try_set_security_ids(token_ids_enum)
    if not tokens_set:
        LOGGER.warning("Failed to configure user token policies (auth=%s); using asyncua defaults", auth)
    server.allow_anonymous = not auth
    try:
        if auth:
            server.supported_tokens = (ua.UserNameIdentityToken,)
        else:
            server.supported_tokens = (ua.AnonymousIdentityToken, ua.UserNameIdentityToken)
    except Exception as exc:
        LOGGER.debug("Could not set server.supported_tokens: %s", exc)
    try:
        if hasattr(server, "iserver") and server.iserver:
            if auth:
                server.iserver.supported_tokens = (ua.UserNameIdentityToken,)
            else:
                server.iserver.supported_tokens = (ua.AnonymousIdentityToken, ua.UserNameIdentityToken)
    except Exception as exc:
        LOGGER.debug("Could not set server.iserver.supported_tokens: %s", exc)

    # Register the user manager only when auth is requested. asyncua 1.1.x
    # calls `user_manager.get_user(iserver, username, password, certificate)`
    # on session activation and treats `None` as reject. For anonymous
    # servers (the default) we leave asyncua's PermissiveUserManager in
    # place so older asyncua versions without the User/UserRole types
    # (e.g. 1.1.0 on platform25) keep working.
    if auth:
        server.iserver.set_user_manager(HardcodedUserManager())
    # Some asyncua versions also read endpoint-level token policies from
    # the server object itself; set it for completeness so endpoints
    # advertise the expected identity tokens.
    if hasattr(server, "user_token_policy"):
        try:
            server.user_token_policy = token_policies
        except Exception as exc:
            LOGGER.debug("server.user_token_policy assign failed: %s", exc)

    try:
        eps = await server.get_endpoints()
        if not eps:
            LOGGER.info("get_endpoints returned no endpoints; user token policies may not be exposed")
        for ep in eps:
            tokens = [
                f"{getattr(t, 'TokenType', '?')}:{getattr(t, 'PolicyId', '?')}:{getattr(t, 'SecurityPolicyUri', '?')}"
                for t in getattr(ep, "UserIdentityTokens", [])
            ]
            LOGGER.info("Configured endpoint %s tokens=%s", getattr(ep, "EndpointUrl", "?"), tokens)
    except Exception as exc:
        LOGGER.debug("Could not read endpoints for logging: %s", exc)
    LOGGER.info(
        "Auth config: auth=%s allow_anonymous=%s user_tokens=%s",
        auth,
        getattr(server, "allow_anonymous", None),
        getattr(server, "user_policy_ids", token_ids_enum),
    )

    uri = "urn:ifw:uatools:pysrv"
    with PROFILER.phase("register_namespace"):
        idx = await server.register_namespace(uri)

    with PROFILER.phase("root system object + state vars"):
        root = await server.nodes.objects.add_object(ua.NodeId("system", idx), "system")
        system_ctx = await add_state_vars(root, "system", "system", idx, device_type="system")
        await add_common_methods(root, system_ctx, idx, "system")

    devices: Dict[str, DeviceContext] = {}
    with PROFILER.phase("build_hierarchy (all subsystems)"):
        await build_hierarchy(idx, root, system_ctx, size, devices)

    # Structured-Variable fixture (additive) — a Variable with child Variables,
    # so the "traverse into structured Variables" fix has local test coverage.
    with PROFILER.phase("structured_variable_fixture"):
        await build_structured_variable_fixture(idx, server.nodes.objects)

    return server, devices, idx


async def build_server_from_xml(port: int, auth: bool, log_level: str,
                                host: str, xml_profile: Path):
    """Replace the built-in namespace with one loaded from a NodeSet2 XML
    file. The standard server setup (endpoint, security, identity
    tokens) is reused via the same code path as ``build_server``, but
    the per-device fixture is NOT added — instead the XML is imported
    into the address space.

    Returns ``(server, idx)`` where ``idx`` is the first non-zero
    namespace index registered by the XML (best effort)."""
    # Reuse the security/auth/identity setup by spinning up a Server
    # the same way build_server() does, then bailing out before the
    # built-in device tree gets created.
    configure_logging(log_level)
    server = Server()
    await server.init()
    server.set_endpoint(f"opc.tcp://{host}:{port}")
    server.set_server_name("uatools-pysrv-xml")
    server.set_security_policy([ua.SecurityPolicyType.NoSecurity])
    # auth handling is intentionally minimal in XML mode — anonymous
    # access only. Use the standard path when you need username auth.
    server.allow_anonymous = True

    await server.import_xml(str(xml_profile))
    # Try to pick a representative ns_idx for callers that want one
    # (e.g. logging). Pick the first non-zero index we find in the
    # namespace array.
    idx = 0
    try:
        ns_array = await server.get_namespace_array()
        for i, _uri in enumerate(ns_array):
            if i != 0:
                idx = i
                break
    except Exception:
        pass
    return server, idx


async def main():
    parser = argparse.ArgumentParser(description="UATK asyncua test server")
    parser.add_argument("-i", "--ip", default="0.0.0.0", help="IP address to bind (default: 0.0.0.0 - all interfaces)")
    parser.add_argument("-p", "--port", type=int, default=4840, help="Server port (default: 4840)")
    parser.add_argument("-a", "--authentication", action="store_true", help="Enable user authentication (user/user_pswd)")
    parser.add_argument(
        "-s", "--size",
        default=SIZE_SMALL,
        choices=["s", "S", "small", "m", "M", "medium", "l", "L", "large", "h", "H", "huge"],
        help="Namespace size: s/S/small (5 subsys, ~4K nodes), m/M/medium (10 subsys, ~15K), l/L/large (10 subsys, ~100K), h/H/huge (10 subsys, ~200K) (default: s)"
    )
    parser.add_argument("-k", "--list-keys", action="store_true", help="List namespace tree and exit")
    parser.add_argument(
        "-l",
        "--log-level",
        default="WARNING",
        help="Log level (ERROR, WARNING, INFO, DEBUG) or logger-specific levels, e.g. asyncua:DEBUG,uatools.pysrv:DEBUG",
    )
    parser.add_argument("--list-loggers", action="store_true", help="List known loggers and exit")
    parser.add_argument(
        "--xml-profile",
        type=Path,
        default=None,
        help="Replace the built-in namespace with one loaded from a NodeSet2 XML "
             "file. The built-in device fixture and simulation loop are NOT used "
             "in this mode.",
    )
    parser.add_argument(
        "--xml-profile-lib",
        default=None,
        help="Importable Python module providing on_<MethodName>(parent, *args) and "
             "on_write_<VariableName>(node, value) callbacks for the loaded XML "
             "profile. The module must be on PYTHONPATH (typically installed under "
             "INTROOT/PREFIX). Without this, the XML namespace is served read-only "
             "with no method behaviour.",
    )
    parser.add_argument(
        "--gen-profile-lib",
        type=Path,
        default=None,
        help="Generate a stub Python module with on_<Method> / on_write_<Variable> "
             "callbacks for every Method and writable Variable in the XML, then "
             "exit. Requires --xml-profile. Refuses to overwrite an existing file "
             "(delete it manually to regenerate). Name mirrors --xml-profile-lib: "
             "the file generated here is the same kind of module you'd later pass "
             "via --xml-profile-lib at run time.",
    )
    parser.add_argument(
        "--profile-startup",
        action="store_true",
        help="Time each startup phase (server.init, namespace build, per-builder "
             "device construction) and print a breakdown once the server is up. "
             "Diagnostic; off by default. No effect on running behaviour.",
    )
    args = parser.parse_args()

    if args.authentication and not AUTH_SUPPORTED:
        configure_logging(args.log_level)
        LOGGER.error(
            "-a / --authentication requires asyncua >= 1.1.4 "
            "(User / UserRole missing from "
            "asyncua.crypto.permission_rules in this version). "
            "Anonymous mode (without -a) still works."
        )
        raise SystemExit(2)

    if args.profile_startup:
        PROFILER.enabled = True

    if args.list_loggers:
        for name in sorted(logging.root.manager.loggerDict.keys()):
            print(name)
        return

    # --- gen-profile-lib: write a stub, then exit. ---------------------
    if args.gen_profile_lib is not None:
        if args.xml_profile is None:
            LOGGER.error("--gen-profile-lib requires --xml-profile.")
            raise SystemExit(2)
        await generate_profile_logic(args.xml_profile, args.gen_profile_lib)
        return

    # --- XML profile mode: replace built-in namespace + skip sim loop. -
    if args.xml_profile is not None:
        configure_logging(args.log_level)
        if not args.xml_profile.exists():
            LOGGER.error("XML profile not found: %s", args.xml_profile)
            raise SystemExit(2)
        server, ns_idx = await build_server_from_xml(
            args.port, args.authentication, args.log_level,
            host=args.ip, xml_profile=args.xml_profile,
        )
        profile_lib = None
        if args.xml_profile_lib:
            try:
                profile_lib = _import_profile_lib(args.xml_profile_lib)
            except Exception as exc:
                LOGGER.error("Could not import --xml-profile-lib %r: %s",
                             args.xml_profile_lib, exc)
                raise SystemExit(2)
        if args.list_keys:
            await list_namespace_tree(server, ns_idx)
            return
        async with server:
            LOGGER.info(
                "Server running on opc.tcp://%s:%s (XML profile=%s, lib=%s)",
                args.ip, args.port, args.xml_profile,
                args.xml_profile_lib or "<none>",
            )
            if profile_lib is not None:
                try:
                    await _bind_profile_logic(server, profile_lib)
                except Exception as exc:
                    LOGGER.warning("Failed to bind profile logic: %s", exc)
            try:
                while True:
                    await asyncio.sleep(3600)
            except asyncio.CancelledError:
                pass
        return

    # --- Default: built-in namespace + simulation loop. ----------------
    size_norm = parse_size(args.size)
    # Override endpoint before init by setting endpoint in Server via helper.
    # asyncua allows set_endpoint after init, so we set it inside build_server using a temporary config.
    t_build_start = time.perf_counter()
    server, devices, ns_idx = await build_server(args.port, args.authentication, args.log_level, host=args.ip, size=size_norm)
    build_secs = time.perf_counter() - t_build_start

    if PROFILER.enabled:
        print(f"\nbuild_server() wall time: {build_secs:.3f}s ({len(devices)} devices)",
              flush=True)
        PROFILER.report()

    if args.list_keys:
        await list_namespace_tree(server, ns_idx)
        return
    sim_task = None
    async with server:
        LOGGER.info("Server running on opc.tcp://%s:%s (auth=%s, size=%s)", args.ip, args.port, args.authentication, size_norm)
        sim_task = asyncio.create_task(simulation_loop(devices))
        try:
            while True:
                await asyncio.sleep(3600)
        except asyncio.CancelledError:
            pass
        finally:
            if sim_task:
                sim_task.cancel()
                with contextlib.suppress(asyncio.CancelledError):
                    await sim_task


if __name__ == "__main__":
    import contextlib

    asyncio.run(main())
