Fix Orca MQTT: reassemble packets across WebSocket messages

OrcaSonar's /mqtt WebSocket proxy writes each ~8 KiB TCP read as its own
WebSocket binary message, so an MQTT packet larger than one read (e.g. a
files.list reply) spans several messages, and several packets can arrive in
one. OrcaMqttConnection assumed exactly one packet per message, so packets
that spanned messages were dropped and the request timed out ("Load failed").

Keep a per-connection byte buffer and extract complete packets by their
remaining-length header, dispatching each. A malformed header (or a declared
size past a 16 MiB cap) drops the connection so the worker reconnects rather
than buffering forever.
This commit is contained in:
Lam Wei Lun
2026-10-07 18:35:12 +08:00
parent 432d6be518
commit 3c646f305c
4 changed files with 302 additions and 5 deletions
+85 -4
View File
@@ -28,6 +28,7 @@
#include <cstdint>
#include <mutex>
#include <cstddef>
#include <boost/system/error_code.hpp>
#include <boost/beast/websocket/rfc6455.hpp>
#include <boost/beast/websocket/stream_base.hpp>
#include <ios>
@@ -53,6 +54,10 @@ struct OrcaMqttConnection::Connection {
boost::asio::ip::tcp::resolver resolver;
boost::asio::steady_timer keepalive_timer;
boost::beast::flat_buffer read_buffer;
// Reassembly buffer for the MQTT byte stream: each WebSocket message carries an
// arbitrary byte range of that stream, so partial packets accumulate here until
// complete. Only the MQTT worker thread touches it.
std::string mqtt_rx;
std::deque<std::shared_ptr<std::vector<uint8_t>>> outbound_packets;
boost::system::error_code terminal_error;
std::atomic_bool async_session_started{false};
@@ -81,6 +86,37 @@ template<class Conn> void expires_never(Conn& conn) {
if (conn.wss) boost::beast::get_lowest_layer(*conn.wss).expires_never();
else if (conn.ws) boost::beast::get_lowest_layer(*conn.ws).expires_never();
}
// Total byte length of the MQTT packet starting at `offset` in `stream`. Returns
// kMqttPacketIncomplete when the fixed header's remaining-length varint is not yet
// complete, or kMqttPacketMalformed when it cannot be valid (more than four length
// bytes, or a declared size past the sanity cap). The caller compares a real length
// against the bytes available.
constexpr std::size_t kMqttPacketIncomplete = 0;
constexpr std::size_t kMqttPacketMalformed = static_cast<std::size_t>(-1);
// MQTT permits ~256 MiB; nothing this client receives is close, so a larger
// declared length is treated as malformed rather than buffered.
constexpr std::size_t kMqttMaxPacketBytes = 16 * 1024 * 1024;
std::size_t mqtt_packet_length(const std::string& stream, std::size_t offset) {
if (offset >= stream.size())
return kMqttPacketIncomplete;
std::size_t multiplier = 1;
std::size_t remaining = 0;
std::size_t i = offset + 1;
for (int len_bytes = 0; len_bytes < 4; ++len_bytes) {
if (i >= stream.size())
return kMqttPacketIncomplete; // varint not complete yet
const std::uint8_t byte = static_cast<std::uint8_t>(stream[i++]);
remaining += static_cast<std::size_t>(byte & 0x7f) * multiplier;
if ((byte & 0x80) == 0) {
const std::size_t total = (i - offset) + remaining;
return total > kMqttMaxPacketBytes ? kMqttPacketMalformed : total;
}
multiplier *= 128;
}
return kMqttPacketMalformed; // more than four length bytes
}
} // namespace
OrcaMqttConnection::~OrcaMqttConnection() { stop(); }
@@ -424,9 +460,9 @@ void OrcaMqttConnection::start_async_read(const std::shared_ptr<Connection>& con
conn->io_context.stop();
return;
}
const std::string packet = boost::beast::buffers_to_string(conn->read_buffer.data());
const std::string message = boost::beast::buffers_to_string(conn->read_buffer.data());
conn->read_buffer.consume(conn->read_buffer.size());
handle_packet(packet);
feed_mqtt(*conn, message);
start_async_read(conn);
};
if (conn->wss)
@@ -435,6 +471,44 @@ void OrcaMqttConnection::start_async_read(const std::shared_ptr<Connection>& con
conn->ws->async_read(conn->read_buffer, std::move(on_read));
}
bool OrcaMqttConnection::drain_mqtt_packets(std::string& stream, std::vector<std::string>& packets) {
std::size_t consumed = 0;
bool ok = true;
while (consumed < stream.size()) {
const std::size_t total = mqtt_packet_length(stream, consumed);
if (total == kMqttPacketIncomplete)
break; // wait for the rest of the stream
if (total == kMqttPacketMalformed) {
ok = false;
break;
}
if (consumed + total > stream.size())
break; // header complete, payload still arriving
packets.emplace_back(stream, consumed, total);
consumed += total;
}
if (consumed > 0)
stream.erase(0, consumed);
if (!ok)
stream.clear(); // drop the poisoned tail; the caller closes the connection
return ok;
}
void OrcaMqttConnection::feed_mqtt(Connection& conn, const std::string& bytes) {
conn.mqtt_rx.append(bytes);
std::vector<std::string> packets;
if (!drain_mqtt_packets(conn.mqtt_rx, packets)) {
// A malformed header can never resync; drop the session so the worker
// reconnects with a clean MQTT stream instead of buffering forever.
if (!stopping.load())
conn.terminal_error = boost::system::errc::make_error_code(boost::system::errc::protocol_error);
conn.io_context.stop();
return;
}
for (const std::string& packet : packets)
handle_packet(packet);
}
void OrcaMqttConnection::schedule_keepalive(const std::shared_ptr<Connection>& conn) {
const int keepalive = current_config.keepalive_seconds;
if (!conn || keepalive <= 0 || stopping.load())
@@ -604,8 +678,9 @@ void OrcaMqttConnection::connect_and_read() {
// disable it before starting the long-lived async WebSocket session.
expires_never(*connection);
const std::string connack = boost::beast::buffers_to_string(buffer.data());
// rc: 0 accepted, 1..5 refusal, -1 malformed/not a CONNACK.
const int rc = (connack.size() == 4 && static_cast<uint8_t>(connack[0]) == 0x20)
// rc: 0 accepted, 1..5 refusal, -1 not a CONNACK. A WebSocket message may
// carry more than the CONNACK; only its first four bytes are the CONNACK.
const int rc = (connack.size() >= 4 && static_cast<uint8_t>(connack[0]) == 0x20)
? static_cast<int>(static_cast<uint8_t>(connack[3]))
: -1;
m_last_connack_rc.store(rc);
@@ -624,6 +699,11 @@ void OrcaMqttConnection::connect_and_read() {
throw std::runtime_error("Orca MQTT CONNECT refused rc=" + std::to_string(rc));
}
// Any bytes the broker sent after the 4-byte CONNACK in the same message
// start the MQTT stream; keep them for the reassembler.
if (connack.size() > 4)
connection->mqtt_rx.assign(connack, 4, std::string::npos);
// The subscription acknowledgement belongs to this MQTT session. Clear
// the previous session's state before notifying the owner, because the
// reconnect callback immediately queues the printer's initial requests.
@@ -636,6 +716,7 @@ void OrcaMqttConnection::connect_and_read() {
notify_state(true);
reconnect_delay_seconds.store(1); // a fresh CONNACK resets the backoff
send_current_subscriptions(connection);
feed_mqtt(*connection, {}); // drain anything that rode in with the CONNACK
start_async_read(connection);
schedule_keepalive(connection);
const std::size_t handlers_run = connection->io_context.run();
+13
View File
@@ -70,6 +70,15 @@ public:
static std::vector<uint8_t> make_subscribe_packet(uint16_t packet_id, const std::string& topic, uint8_t qos);
static std::vector<uint8_t> make_unsubscribe_packet(uint16_t packet_id, const std::string& topic);
// Split a raw MQTT byte stream into complete packets: each complete packet is
// appended to `packets` and erased from `stream`; a partial trailing packet is
// left in `stream` for the next call. Returns false when the stream begins with
// a malformed header, which cannot resync — the caller must drop the connection.
// Public for unit tests. MQTT-over-WebSocket allows a packet to span frames and
// several packets per frame, so the receiver must reassemble the stream rather
// than assume one packet per WebSocket message.
static bool drain_mqtt_packets(std::string& stream, std::vector<std::string>& packets);
~OrcaMqttConnection();
bool start(const Config& config, MessageHandler on_message, StateHandler on_state);
@@ -118,6 +127,10 @@ private:
void send_current_subscriptions(const std::shared_ptr<Connection>& conn);
void send_pending_subscriptions(const std::shared_ptr<Connection>& conn);
void handle_packet(const std::string& packet);
// Append freshly received WebSocket bytes to `conn`'s MQTT stream and dispatch
// every complete packet. An MQTT packet may span several WebSocket messages and
// several packets may arrive in one, so the stream is reassembled here.
void feed_mqtt(Connection& conn, const std::string& bytes);
void notify_state(bool is_now_connected);
void run();
+22 -1
View File
@@ -77,7 +77,11 @@ class MockBroker
public:
// refuse_auth: answer every CONNECT with CONNACK rc 5 (not authorized) and
// close, so the reconnect/refusal paths can be exercised.
explicit MockBroker(bool refuse_auth = false) : m_refuse_auth(refuse_auth), m_acceptor(m_io)
// stall_connack: accept the WebSocket upgrade and the CONNECT but never answer,
// so the client's synchronous handshake can be exercised against a peer that
// stops responding.
explicit MockBroker(bool refuse_auth = false, bool stall_connack = false)
: m_refuse_auth(refuse_auth), m_stall_connack(stall_connack), m_acceptor(m_io)
{
const tcp::endpoint endpoint(net::ip::make_address("127.0.0.1"), 0);
m_acceptor.open(endpoint.protocol());
@@ -124,6 +128,20 @@ public:
m_stream->write(net::buffer(packet), ec); // a vanished client is not a test failure
}
// Write arbitrary bytes as one binary WebSocket message, without the topic
// filter push_report applies. Lets a test split a single MQTT packet across
// messages or coalesce several packets into one - the two framings the real
// OrcaSonar /mqtt proxy produces with its 8 KiB reads.
void push_raw(const std::string& bytes)
{
std::lock_guard<std::mutex> lock(m_mutex);
if (!m_stream || !m_stream_ready)
return;
boost::system::error_code ec;
m_stream->binary(true);
m_stream->write(net::buffer(bytes), ec);
}
// Force-close the live client socket; the worker's read returns an error and
// the accept loop picks up the client's reconnect.
void drop_client()
@@ -237,6 +255,8 @@ private:
write_packet(stream, {0x20, 0x02, 0x00, 0x05}); // CONNACK not authorized
return false;
}
if (m_stall_connack)
return true; // accepted, never answered: the client's read blocks
write_packet(stream, {0x20, 0x02, 0x00, 0x00}); // CONNACK accepted
return true;
}
@@ -364,6 +384,7 @@ private:
}
const bool m_refuse_auth;
const bool m_stall_connack;
net::io_context m_io;
tcp::acceptor m_acceptor;
std::string m_port;
@@ -252,3 +252,185 @@ TEST_CASE("OrcaMqtt auth rejection is terminal (no retry storm)", "[OrcaMqtt]")
CHECK(broker.connect_count() <= 2);
conn.stop();
}
// --- MQTT-over-WebSocket stream reassembly. The OrcaSonar /mqtt proxy writes each
// 8 KiB TCP read as its own WebSocket message, so a packet larger than one read
// (e.g. a files.list reply) spans several messages; and several small packets can
// arrive in one. The receiver must reassemble the MQTT stream, not assume one
// packet per message.
TEST_CASE("OrcaMqtt drain_mqtt_packets reassembles across message boundaries", "[OrcaMqtt]") {
const std::vector<std::uint8_t> a = OrcaMqttConnection::make_publish_packet("device/a/report", R"({"n":1})");
const std::vector<std::uint8_t> b = OrcaMqttConnection::make_publish_packet("device/a/report", R"({"n":2})");
const std::string pa(a.begin(), a.end());
const std::string pb(b.begin(), b.end());
// A partial packet yields nothing and is retained verbatim.
std::string stream = pa.substr(0, pa.size() / 2);
std::vector<std::string> packets;
OrcaMqttConnection::drain_mqtt_packets(stream, packets);
CHECK(packets.empty());
CHECK(stream == pa.substr(0, pa.size() / 2));
// Completing the packet yields it and drains the stream.
stream += pa.substr(pa.size() / 2);
OrcaMqttConnection::drain_mqtt_packets(stream, packets);
REQUIRE(packets.size() == 1);
CHECK(packets[0] == pa);
CHECK(stream.empty());
// Two packets coalesced in one buffer are both extracted, in order.
stream = pa + pb;
packets.clear();
OrcaMqttConnection::drain_mqtt_packets(stream, packets);
REQUIRE(packets.size() == 2);
CHECK(packets[0] == pa);
CHECK(packets[1] == pb);
CHECK(stream.empty());
// A trailing partial packet is kept for the next call.
stream = pa + pb.substr(0, 2);
packets.clear();
CHECK(OrcaMqttConnection::drain_mqtt_packets(stream, packets));
REQUIRE(packets.size() == 1);
CHECK(packets[0] == pa);
CHECK(stream == pb.substr(0, 2));
// A multibyte remaining length (payload > 127 bytes) is parsed correctly.
const std::string big(300, 'x');
const auto big_packet = OrcaMqttConnection::make_publish_packet("device/a/report", big);
const std::string pb_big(big_packet.begin(), big_packet.end());
stream = pb_big;
packets.clear();
REQUIRE(OrcaMqttConnection::drain_mqtt_packets(stream, packets));
REQUIRE(packets.size() == 1);
CHECK(packets[0] == pb_big);
CHECK(stream.empty());
// A malformed header is reported so the caller can drop the connection.
const std::vector<std::uint8_t> malformed_bytes{0x30, 0x80, 0x80, 0x80, 0x80};
stream.assign(malformed_bytes.begin(), malformed_bytes.end());
packets.clear();
CHECK_FALSE(OrcaMqttConnection::drain_mqtt_packets(stream, packets));
CHECK(packets.empty());
CHECK(stream.empty());
}
TEST_CASE("OrcaMqtt delivers a report split across WebSocket messages", "[OrcaMqtt]") {
orca_mqtt_test::MockBroker broker;
// Callback state is declared before `conn` so it outlives the worker: if a
// REQUIRE fails, `conn`'s destructor still runs before this state is destroyed.
std::mutex m;
std::condition_variable cv;
std::vector<std::string> got;
OrcaMqttConnection conn;
OrcaMqttConnection::Config cfg;
cfg.url = broker.ws_url();
cfg.use_tls = false;
cfg.username = "u";
cfg.password = "p";
REQUIRE(conn.start(cfg,
[&](const std::string&, const std::string& p) {
{ std::lock_guard<std::mutex> l(m); got.push_back(p); }
cv.notify_all();
},
[](bool, bool) {}));
REQUIRE(conn.subscribe("dev-1"));
REQUIRE(wait_subscribed(broker, "device/dev-1/report"));
const std::string payload = R"({"print":{"command":"push_status","sequence_id":"30001","result":"success"}})";
const auto packet = OrcaMqttConnection::make_publish_packet("device/dev-1/report", payload);
const std::string bytes(packet.begin(), packet.end());
const std::size_t third = bytes.size() / 3;
broker.push_raw(bytes.substr(0, third));
broker.push_raw(bytes.substr(third, third));
broker.push_raw(bytes.substr(2 * third));
bool delivered = false;
{
std::unique_lock<std::mutex> l(m);
delivered = cv.wait_for(l, std::chrono::seconds(3), [&] { return !got.empty(); });
}
REQUIRE(delivered);
{
std::lock_guard<std::mutex> l(m);
REQUIRE(got.size() == 1);
CHECK(got.front().find("push_status") != std::string::npos);
}
conn.stop();
}
TEST_CASE("OrcaMqtt delivers two reports coalesced into one WebSocket message", "[OrcaMqtt]") {
orca_mqtt_test::MockBroker broker;
// Callback state is declared before `conn` so it outlives the worker: if a
// REQUIRE fails, `conn`'s destructor still runs before this state is destroyed.
std::mutex m;
std::condition_variable cv;
std::vector<std::string> got;
OrcaMqttConnection conn;
OrcaMqttConnection::Config cfg;
cfg.url = broker.ws_url();
cfg.use_tls = false;
cfg.username = "u";
cfg.password = "p";
REQUIRE(conn.start(cfg,
[&](const std::string&, const std::string& p) {
{ std::lock_guard<std::mutex> l(m); got.push_back(p); }
cv.notify_all();
},
[](bool, bool) {}));
REQUIRE(conn.subscribe("dev-1"));
REQUIRE(wait_subscribed(broker, "device/dev-1/report"));
const auto p1 = OrcaMqttConnection::make_publish_packet("device/dev-1/report", R"({"print":{"command":"push_status","sequence_id":"31001"}})");
const auto p2 = OrcaMqttConnection::make_publish_packet("device/dev-1/report", R"({"print":{"command":"push_status","sequence_id":"31002"}})");
const std::string both(std::string(p1.begin(), p1.end()) + std::string(p2.begin(), p2.end()));
broker.push_raw(both);
bool two = false;
{
std::unique_lock<std::mutex> l(m);
two = cv.wait_for(l, std::chrono::seconds(3), [&] { return got.size() >= 2; });
}
REQUIRE(two);
{
std::lock_guard<std::mutex> l(m);
CHECK(got[0].find("31001") != std::string::npos);
CHECK(got[1].find("31002") != std::string::npos);
}
conn.stop();
}
TEST_CASE("OrcaMqtt stop() unblocks a stalled handshake", "[OrcaMqtt]") {
// A peer that completes the WebSocket upgrade and CONNECT but never answers:
// the worker blocks in the synchronous CONNACK read. stop() must shut the
// socket down and join it promptly rather than wait out the 10s deadline.
orca_mqtt_test::MockBroker broker(/*refuse_auth=*/false, /*stall_connack=*/true);
OrcaMqttConnection conn;
OrcaMqttConnection::Config cfg;
cfg.url = broker.ws_url();
cfg.use_tls = false;
cfg.username = "u";
cfg.password = "p";
// start() waits up to 10s for CONNACK, so run it off the test thread.
std::thread starter([&] { conn.start(cfg, [](const std::string&, const std::string&) {}, [](bool, bool) {}); });
struct FinalJoin {
std::thread& thread;
~FinalJoin() {
if (thread.joinable())
thread.join();
}
} final_join{starter};
// Wait until the broker has taken the CONNECT: the worker is now in the read.
for (int i = 0; i < 300 && broker.connect_count() == 0; ++i)
std::this_thread::sleep_for(std::chrono::milliseconds(10));
REQUIRE(broker.connect_count() > 0);
const auto started_at = std::chrono::steady_clock::now();
conn.stop();
const auto elapsed = std::chrono::steady_clock::now() - started_at;
CHECK(elapsed < std::chrono::seconds(3));
CHECK_FALSE(conn.is_running());
}