/*! * * Copyright (c) 2026 Diality Inc. - All Rights Reserved. * \copyright * THIS CODE MAY NOT BE COPIED OR REPRODUCED IN ANY FORM, IN PART OR IN * WHOLE, WITHOUT THE EXPLICIT PERMISSION OF THE COPYRIGHT OWNER. * * \file MqttTcpTransport.cpp * \author (original) Stephen Quong * \date (original) 30-Jul-2026 * */ #include #include "MqttTcpTransport.h" Q_LOGGING_CATEGORY(lcMqttTcp, "cloudconnect.mqtt.tcp") namespace { constexpr int KEEPALIVE_DIVISOR = 2; } /*! * \brief MqttTcpTransport::MqttTcpTransport * \details Constructor * \param parent - optional QObject parent */ MqttTcpTransport::MqttTcpTransport(QObject *parent) : QObject(parent) { connect(&_socket, &QTcpSocket::connected, this, &MqttTcpTransport::onSocketConnected); connect(&_socket, &QTcpSocket::disconnected, this, &MqttTcpTransport::onSocketDisconnected); connect(&_socket, &QTcpSocket::readyRead, this, &MqttTcpTransport::onReadyRead); connect(&_socket, &QAbstractSocket::errorOccurred, this, [this](QAbstractSocket::SocketError error) { qCWarning(lcMqttTcp).noquote() << "Socket error:" << _socket.errorString() << QString("(%1)").arg(static_cast(error)); setSessionState(false); }); _keepAliveTimer.setSingleShot(false); connect(&_keepAliveTimer, &QTimer::timeout, this, &MqttTcpTransport::onKeepAliveTimer); } /*! * \brief MqttTcpTransport::~MqttTcpTransport * \details Drops the callbacks before closing so teardown cannot re-enter the owner. */ MqttTcpTransport::~MqttTcpTransport() { clearCallbacks(); close(); } /*! * \brief MqttTcpTransport::open * \details Starts the TCP connection. CONNECT is sent once the socket is up, and * the session is not established until CONNACK arrives. * \param config - endpoint, port, client id and keep-alive * \return true if the connection attempt was started. */ bool MqttTcpTransport::open(const Config &config) { _config = config; if (_socket.state() != QAbstractSocket::UnconnectedState) { qCWarning(lcMqttTcp) << "open() while the socket is already in use"; return false; } qCInfo(lcMqttTcp).noquote() << "Connecting to" << _config.endpoint << "port" << _config.port; _socket.connectToHost(_config.endpoint, _config.port); return true; } /*! * \brief MqttTcpTransport::close * \details Sends DISCONNECT when a session is up, then closes the socket. */ void MqttTcpTransport::close() { if (_sessionUp && _socket.state() == QAbstractSocket::ConnectedState) { _socket.write(Mqtt::buildDisconnect()); _socket.flush(); } _keepAliveTimer.stop(); if (_socket.state() != QAbstractSocket::UnconnectedState) { _socket.disconnectFromHost(); } else { setSessionState(false); } } /*! * \brief MqttTcpTransport::isConnected * \details Reports whether the MQTT session is established. * \return true once CONNACK has been accepted and before the link drops. */ bool MqttTcpTransport::isConnected() const { return _sessionUp; } /*! * \brief MqttTcpTransport::publish * \details Sends one QoS 1 PUBLISH and records the packet identifier so the * matching PUBACK can be tied back to the caller's message. * \param topic - fully resolved topic name * \param payload - message bytes, sent unmodified * \param messageId - caller token to report on acknowledgement, or -1 * \return true if the packet was written to the socket. */ bool MqttTcpTransport::publish(const QString &topic, const QByteArray &payload, qint64 messageId) { if (!_sessionUp) { return false; } const quint16 packetId = nextPacketId(); _pending.insert(packetId, messageId); const QByteArray packet = Mqtt::buildPublish(topic, payload, 1, packetId); if (_socket.write(packet) != packet.size()) { qCWarning(lcMqttTcp).noquote() << "Short write publishing to" << topic; _pending.remove(packetId); return false; } return true; } /*! * \brief MqttTcpTransport::setCallbacks * \details Installs the state and acknowledgement callbacks. * \param onState - invoked on session establishment and loss * \param onAck - invoked once per publish outcome */ void MqttTcpTransport::setCallbacks(StateCallback onState, AckCallback onAck) { _onState = std::move(onState); _onAck = std::move(onAck); } /*! * \brief MqttTcpTransport::clearCallbacks * \details Drops both callbacks. */ void MqttTcpTransport::clearCallbacks() { _onState = nullptr; _onAck = nullptr; } /*! * \brief MqttTcpTransport::onSocketConnected * \details Sends CONNECT once the TCP connection is up. */ void MqttTcpTransport::onSocketConnected() { qCInfo(lcMqttTcp).noquote() << "Socket connected, sending CONNECT as" << _config.clientId; _socket.write(Mqtt::buildConnect(_config.clientId, _config.keepAliveSecs, _config.cleanSession)); } /*! * \brief MqttTcpTransport::onSocketDisconnected * \details Tears down session state when the link drops. */ void MqttTcpTransport::onSocketDisconnected() { qCInfo(lcMqttTcp).noquote() << "Socket disconnected from" << _config.endpoint; setSessionState(false); } /*! * \brief MqttTcpTransport::onReadyRead * \details Accumulates inbound bytes and drains whole packets out of the buffer. */ void MqttTcpTransport::onReadyRead() { _rxBuf.append(_socket.readAll()); Mqtt::ReadState state = Mqtt::ReadState::Incomplete; do { Mqtt::Packet packet; state = Mqtt::readPacket(_rxBuf, packet); if (state == Mqtt::ReadState::Complete) { handlePacket(packet); } else if (state == Mqtt::ReadState::Error) { qCCritical(lcMqttTcp) << "Malformed packet, dropping the connection"; _rxBuf.clear(); _socket.abort(); setSessionState(false); } } while (state == Mqtt::ReadState::Complete && !_rxBuf.isEmpty()); } /*! * \brief MqttTcpTransport::onKeepAliveTimer * \details Sends PINGREQ so the broker does not time the session out. */ void MqttTcpTransport::onKeepAliveTimer() { if (_sessionUp) { _socket.write(Mqtt::buildPingReq()); } } /*! * \brief MqttTcpTransport::handlePacket * \details Dispatches one decoded inbound packet. * \param packet - the decoded packet */ void MqttTcpTransport::handlePacket(const Mqtt::Packet &packet) { switch (packet.type) { case Mqtt::PacketType::ConnAck: if (packet.returnCode == Mqtt::ConnAckCode::Accepted) { qCInfo(lcMqttTcp).noquote() << "CONNACK accepted by" << _config.endpoint; if (_config.keepAliveSecs > 0) { _keepAliveTimer.start(_config.keepAliveSecs * 1000 / KEEPALIVE_DIVISOR); } setSessionState(true); } else { qCCritical(lcMqttTcp).noquote() << "CONNACK refused with code" << static_cast(packet.returnCode); setSessionState(false); _socket.disconnectFromHost(); } break; case Mqtt::PacketType::PubAck: { // take() returns a default-constructed value for an unknown key, so an // unsolicited PUBACK is reported against -1 rather than silently dropped. if (!_pending.contains(packet.packetId)) { qCWarning(lcMqttTcp) << "PUBACK for an unknown packet id" << packet.packetId; break; } const qint64 messageId = _pending.take(packet.packetId); if (_onAck) { _onAck(messageId, true); } break; } case Mqtt::PacketType::PingResp: break; default: qCWarning(lcMqttTcp).noquote() << "Unexpected" << Mqtt::typeName(packet.type) << "from the broker"; break; } } /*! * \brief MqttTcpTransport::setSessionState * \details Applies a session transition and reports it, suppressing repeats. * \param established - true when the session came up */ void MqttTcpTransport::setSessionState(bool established) { if (_sessionUp == established) { return; } _sessionUp = established; if (!established) { _keepAliveTimer.stop(); _rxBuf.clear(); failPending(); } if (_onState) { _onState(established); } } /*! * \brief MqttTcpTransport::failPending * \details Reports every unacknowledged publish as failed. */ void MqttTcpTransport::failPending() { if (_pending.isEmpty()) { return; } qCWarning(lcMqttTcp) << "Failing" << _pending.size() << "unacknowledged publishes"; const QHash pending = _pending; _pending.clear(); if (_onAck) { for (auto it = pending.constBegin(); it != pending.constEnd(); ++it) { _onAck(it.value(), false); } } } /*! * \brief MqttTcpTransport::nextPacketId * \details Produces the next packet identifier. * \return A value in 1..65535. */ quint16 MqttTcpTransport::nextPacketId() { ++_lastPacketId; if (_lastPacketId == 0) { _lastPacketId = 1; } return _lastPacketId; }