diff --git a/src/slic3r/Utils/OrcaMqttConnection.cpp b/src/slic3r/Utils/OrcaMqttConnection.cpp index 8bd3603b93..e475afd6ee 100644 --- a/src/slic3r/Utils/OrcaMqttConnection.cpp +++ b/src/slic3r/Utils/OrcaMqttConnection.cpp @@ -28,6 +28,7 @@ #include #include #include +#include #include #include #include @@ -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>> outbound_packets; boost::system::error_code terminal_error; std::atomic_bool async_session_started{false}; @@ -81,6 +86,37 @@ template 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(-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(stream[i++]); + remaining += static_cast(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& 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& con conn->ws->async_read(conn->read_buffer, std::move(on_read)); } +bool OrcaMqttConnection::drain_mqtt_packets(std::string& stream, std::vector& 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 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& 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(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(connack[0]) == 0x20) ? static_cast(static_cast(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(); diff --git a/src/slic3r/Utils/OrcaMqttConnection.hpp b/src/slic3r/Utils/OrcaMqttConnection.hpp index ef9593a533..655def72f8 100644 --- a/src/slic3r/Utils/OrcaMqttConnection.hpp +++ b/src/slic3r/Utils/OrcaMqttConnection.hpp @@ -70,6 +70,15 @@ public: static std::vector make_subscribe_packet(uint16_t packet_id, const std::string& topic, uint8_t qos); static std::vector 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& 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& conn); void send_pending_subscriptions(const std::shared_ptr& 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(); diff --git a/tests/slic3rutils/orca_mqtt_mock_broker.hpp b/tests/slic3rutils/orca_mqtt_mock_broker.hpp index 29b90b9fac..ee49f1494c 100644 --- a/tests/slic3rutils/orca_mqtt_mock_broker.hpp +++ b/tests/slic3rutils/orca_mqtt_mock_broker.hpp @@ -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 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; diff --git a/tests/slic3rutils/test_orca_mqtt_connection.cpp b/tests/slic3rutils/test_orca_mqtt_connection.cpp index 130495f3ab..bf77e3e3d3 100644 --- a/tests/slic3rutils/test_orca_mqtt_connection.cpp +++ b/tests/slic3rutils/test_orca_mqtt_connection.cpp @@ -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 a = OrcaMqttConnection::make_publish_packet("device/a/report", R"({"n":1})"); + const std::vector 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 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 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 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 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 l(m); + delivered = cv.wait_for(l, std::chrono::seconds(3), [&] { return !got.empty(); }); + } + REQUIRE(delivered); + { + std::lock_guard 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 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 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 l(m); + two = cv.wait_for(l, std::chrono::seconds(3), [&] { return got.size() >= 2; }); + } + REQUIRE(two); + { + std::lock_guard 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()); +}