MQTT binding#

Carrying values over MQTT topics.

The MQTT semantic connector: the first instance of the binding seam.

A parameter is reachable over MQTT when its domain property is a subproperty of inf:isInterfaceAccessibleMQTTParameter. The parameter node then carries a broker address and a read topic. A settable parameter node also carries a set topic.

Two connectors per settable parameter, one ConnectionInfo. MqttClientConnector takes a single topic. Its consume() publishes to the topic it subscribed to. So it physically cannot serve a read topic plus a distinct set topic. ConnectionRegistry.connections is Dict[ConnectionInfo, List[str]]. So both bind to one ConnectionInfo, and differ only in sync_direction — precisely what SyncDirection is for. For the scenario-3 TransferUnit that is 4 parameters, 4 bindings, 6 connectors, 6 topics.

aiomqtt is an optional extra of transitional_sync_middleware (its industrial group). So the connector class is imported lazily. Importing this module must not require a working MQTT stack. The registry must exist for every connector wiring, including one that wires nothing. An inspecting instance with no aiomqtt installed must still receive the projection.

class kapps_semantic_middleware.connectors.mqtt_binding.MQTTParameterFormatter(model_type: type, northbound_facets: Dict[str, Any], value_field: str, value_path: str | None = None, *, parameter_label: str = '', topic: str = '', set_topic: str = '')[source]#

Bases: object

Translates between an MQTT payload and the persistence value of a parameter node.

The persistence value of a COMPLEX property is a list of one generated model, not a scalar: [AnonymousClass(hasValue=[12.1], hasUnit=['m/s'], accessMode=['readwrite'])]. ConnectionInfo bottoms out at that property. field_id is a plain getattr and the node’s inf:hasValue is one level further down. The formatter bridges the device scalar and the node the middleware persists.

Payload shape. Raw scalar by default. If the parameter declares inf:hasMQTTValuePath, use a JSON envelope with the value at that dotted path. The path is one property and is honoured symmetrically on both read and write.

Symmetry. MqttClientConnector is asymmetric. Its listener runs json.loads(payload) on everything it receives. consume() publishes its argument raw. So deserialize receives an already-parsed value. serialize must produce the encoded bytes.

Why the node is reassembled rather than replaced by a bare value. update_persistence_with_value does setattr(contained_model, field_id, value), replacing the whole list. Formatter.deserialize sees only the payload. It has no access to the current persistence value. A bare scalar would therefore blank the unit and the access mode in the model served over REST. The graph layer now diffs per triple, so unchanged facets are preserved there. The in-memory half remains, and it is what this reassembly is for. The northbound payload keeps its unit after the first device message. The facets come from the same metadata the binding already read. So this is a pure function per message, with no read of current state. _last_value exists solely for change-detection in logging, and never affects the returned value.

deserialize(data: Any) → List[Any][source]#

Device payload to the persistence value (a one-element list holding the node).

serialize(data: Any) → bytes[source]#

Persistence value to the device payload, encoded for consume.

Raises if nothing has ever been observed for this parameter. Under the locator pattern the graph holds no inf:hasValue until a connector fills it. A bare None would encode as the JSON scalar null, not a valid device value, and a receiver expecting a number or boolean chokes on it. The caller’s best-effort fan-out (persisted_connector.py) catches and logs a per-connector failure, so raising here is the same “nothing safe to send” outcome that silently accepting None would have hidden as a network-level failure at the device.

This used to be the ordinary case rather than the guard it reads as. The fan-out asked every write leg on a resource to re-derive its slice whenever any field changed, so a parameter nothing had ever set was routinely asked to publish. The fan-out is now scoped to the field that actually moved (changed), so a write leg is now only asked when its own parameter was written – and a write carries the value with it. The guard stays because “asked to publish something never observed” is still a real error state, and one worth naming rather than encoding as null.

class kapps_semantic_middleware.connectors.mqtt_binding.MQTTBinding[source]#

Bases: object

Binds an MQTT-reachable parameter to one or two MqttClientConnector instances.

static build(binding: ParameterBinding, direction: SyncDirection, *, ensure_transport: Callable[[str, int], None] | None = None) → Iterable[Registration][source]#

One read registration always. One write registration when both sides permit it.

direction has already been reduced to the most restrictive of the parameter’s inf:accessMode and the instance’s flavour, so this only honours it. It must never re-widen.

ensure_transport, when given, is called once with the declared (host, port), immediately before the first connector for it is built. Deduping across several parameters that share one address is the caller’s job (wiring.py’s plan_wiring), not this method’s – a single build call sees only its own parameter, never its siblings.