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:
objectTranslates 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'])].ConnectionInfobottoms out at that property.field_idis a plaingetattrand the node’sinf:hasValueis 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.
MqttClientConnectoris asymmetric. Its listener runsjson.loads(payload)on everything it receives.consume()publishes its argument raw. Sodeserializereceives an already-parsed value.serializemust produce the encoded bytes.Why the node is reassembled rather than replaced by a bare value.
update_persistence_with_valuedoessetattr(contained_model, field_id, value), replacing the whole list.Formatter.deserializesees 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_valueexists 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:hasValueuntil a connector fills it. A bareNonewould encode as the JSON scalarnull, 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 acceptingNonewould 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 asnull.
- class kapps_semantic_middleware.connectors.mqtt_binding.MQTTBinding[source]#
Bases:
objectBinds an MQTT-reachable parameter to one or two
MqttClientConnectorinstances.- 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.
directionhas already been reduced to the most restrictive of the parameter’sinf:accessModeand 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’splan_wiring), not this method’s – a singlebuildcall sees only its own parameter, never its siblings.