Composite Device
CompositeDevice groups several devices — local instances and/or remote
TCP/IP clients — behind one object, so an orchestration (typically an
Experiment) manages a whole setup as a single unit.
Building a composite
The shortest form names the sub-devices and lets the composite open the remote clients itself:
from plesty.lib.device.composite_device import CompositeDevice
composite = CompositeDevice(["hwp", "spec", "pm"])
Server addresses are credentials and stay in .env; the composite derives each
variable from the sub-device’s own attribute name — hwp → HWP_ADDRESS — so
nothing has to be repeated in code or committed to a config file:
# .env
HWP_ADDRESS=tcp://192.168.1.20:5555
SPEC_ADDRESS=tcp://192.168.1.21:5555
PM_ADDRESS=tcp://192.168.1.22:5555
A deployment with its own variable namespace passes env_prefix
(cam → PL_SCAN_CAM_ADDRESS), and an already-loaded EnvSettings can be
injected as env:
composite = CompositeDevice(["cam", "led", "amc"], env_prefix="PL_SCAN")
A sub-device value takes one of four forms, and they mix freely in one mapping:
Value |
Meaning |
|---|---|
Device instance |
A |
|
The ZMQ address of its device server |
|
A configuration (keys below) |
|
An all-default configuration — same as a bare name |
from plesty.lib.device.composite_device import CompositeDevice, connect_client
composite = CompositeDevice(
devices={
"stage": stage_device, # local BaseDeviceSyncModel instance
"spec": {"timeout_ms": 5_000}, # remote; address from SPEC_ADDRESS
"cam": "tcp://192.168.1.20:5555", # remote; explicit address
"pm": {"env": "LEGACY_POWERMETER"}, # remote; non-matching variable name
},
)
Config key |
Purpose |
|---|---|
|
Explicit ZMQ address, skipping the environment |
|
Environment variable to read, overriding the derived name |
|
Forwarded to |
|
Budget for the first connect (default 10 s); past it the constructor raises |
|
Remote methods the server must expose; checked by |
The address comes from address if given, else the environment; an unset
variable raises a ValueError naming the variable to set.
Pass an instance when you need control over how the client is opened — for example building it on its own single-worker thread to keep ZMQ socket affinity. That cannot be expressed as configuration:
composite = CompositeDevice({"cam": connect_client("tcp://192.168.1.20:5555")})
Every sub-device is reconnectable with no address written twice: a configured
one registers the address it resolved, and a client passed as an instance
reports the one it was opened with (client.address).
Sub-devices are reached by name (composite.stage, composite.cam).
connect_all() / disconnect_all() handle every local sub-device, and the
composite is a context manager:
with composite:
print(composite.identity()) # {"stage": ..., "cam": ...}
Remote clients manage their own connection.
Resilient remote calls
Long-running orchestrations must survive transient network hiccups. Instead
of calling a remote sub-device method directly, route it through call():
frame = composite.call("cam", "acquire", exposure_s=0.1)
On a TimeoutError, call() rebuilds the sub-device’s client — discarding
the stale ZMQ reply left in the socket, which would otherwise poison every
subsequent request — waits briefly, and retries (default 3 attempts) before
the final timeout propagates. reconnect(dev) performs the same rebuild
explicitly.
connect_client(address) is the module-level builder used for the rebuilds:
build_client plus a connection-handshake retry with linear backoff, which
matters when several clients open back-to-back against the same server.
reconnect(dev) rebuilds with the sub-device’s own timeout_ms /
attempts / backoff_s, not the defaults.
Every multi-device composite gets this robustness for free instead of re-implementing it per module.
Connection state and a bounded first connect
The composite is never silent about its servers. Constructing it logs one
line per sub-device (Connecting hwp @ tcp://… (timeout 5000 ms) …,
Connected hwp @ tcp://… in 0.4 s, or Connect to tcp://… failed, attempt 1/4: … ; retrying in 2 s) and a summary (PolPlRig ready in 1.2 s: hwp=tcp://…, spec=tcp://…) — the console sink is switched on by the
composite itself (ensure_logging()), so this shows even when the composite
is built before the experiment that would configure logging.
The first connect is bounded: each configured sub-device gets
connect_deadline_s (default 10 s); a server that does not answer within it
raises DeviceUnreachableError(dev, address, elapsed_s) from the constructor
instead of grinding through attempts × timeout × backoff. Under thread
affinity the first connects run in parallel on the sub-devices’ own threads,
so three dead servers cost one deadline, not three. Reconnects during a run
keep the attempt ladder — a server busy finishing a long move must be waited
for, not declared dead.
Each sub-device carries a ConnectionState — CONNECTING → CONNECTED → DEGRADED (a call timed out, the client is being rebuilt) → DISCONNECTED
(connect/reconnect failed, or disconnect_all()):
composite.status() # {"hwp": ConnectionState.CONNECTED, "spec": ...}
composite.status("hwp") # ConnectionState.CONNECTED
@composite.on_status_change
def _show(dev, old, new, detail):
ui.set_led(dev, new) # runs on the thread that changed the state
Every transition is logged (INFO, or WARNING for DEGRADED /
DISCONNECTED) whether or not anyone subscribes; a subscriber that raises
is logged and ignored — a broken UI hook must not break the rig.
One thread per sub-device
ZMQ sockets are not thread-safe, and even migrating one between threads is
fragile — a rig that drove three instruments from a shared pool lost messages
under real network latency until every socket was pinned. thread_affinity
gives each sub-device its own thread, opens its client there, and keeps every
later call on it:
rig = CompositeDevice(["hwp", "spec", "pm"], thread_affinity=True)
submit(dev, func, ...) returns a Future, which is how two instruments cover
the same time window — the exposure and the power reading below start together
rather than one after the other:
spectrum = rig.submit("spec", "acquire", timeout=exposure_s + 5)
power = rig.submit("pm", "measure_power", averaging=10)
row = (spectrum.result(), power.result())
call() hops onto the owning thread by itself, so existing code keeps working;
a call made from a sub-device’s own thread runs inline instead of waiting on
the single worker it already occupies. Without thread_affinity, submit()
runs inline and hands back a resolved future, so the same code stays correct —
just sequential. shutdown() (also called by disconnect_all()) stops the
threads.
Routing a session’s data
set_data_path(path) points every sub-device at one session directory, so a run’s
acquisitions land together instead of in each device’s configured default:
composite.set_data_path("2026-08-11/run-3")
composite.set_data_path("2026-08-11/run-3", subdirs=True) # <path>/cam, <path>/meter
It returns a {sub-device: resolved path} mapping. A sub-device that rejects the path
(outside its allowed roots — see Base Device) aborts the
call, so a run never proceeds believing its data is routed when it is not.
That strictness is also why the routing can be narrowed. In a rig where only one instrument writes data, naming it is the honest form — asking a rotation stage for a data directory either fails or means nothing:
composite.set_data_path(frames_dir, devices=["spec"])
Preflight
preflight() proves each sub-device is the device it is supposed to be and
answers: a remote sub-device must expose every method in its requires
list and its identity() must reply; a local one must have the required
attributes. It returns {name: problem} for the failures (empty when the
rig is ready) and logs one line per sub-device either way. Experiment
runs it in setup(); plesty.lib.experiment.connect_rig runs it at
start-up and asks the operator what to do (see the experiment guide).
Interplay with experiments
Experiment accepts a CompositeDevice at construction; its default
setup()/teardown() call connect_all()/disconnect_all(). Step
operations then drive sub-devices via self.devices.<name> or, for remote
calls that should retry, self.devices.call(...).