mirror of
https://github.com/OrcaSlicer/OrcaSlicer.git
synced 2026-10-07 15:51:08 +00:00
OrcaPrinterAgent - Files over MQTT (#16251)
- Add packet draining to OrcaMqttConnection - Change files list/delete to use the Mqtt version instead of http
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -156,14 +156,15 @@ std::string http_origin_from_lan_ws(const std::string& ws_url)
|
||||
return s;
|
||||
}
|
||||
|
||||
// print.gcode_file is non-idempotent and OrcaSonar replays a cached response for a
|
||||
// reused (namespace, command, sequence_id). Seed from the wall clock so ids do not
|
||||
// collide across slicer restarts, then bump once per call within a run.
|
||||
std::string next_gcode_file_sequence_id()
|
||||
// print.gcode_file and files.* are non-idempotent and OrcaSonar replays a cached
|
||||
// response for a reused (namespace, command, sequence_id). Seed from a
|
||||
// high-resolution wall clock so a fast restart cannot reissue an id a prior run
|
||||
// used, then bump once per call within a run.
|
||||
std::string next_command_sequence_id()
|
||||
{
|
||||
static std::atomic<uint64_t> counter{[] {
|
||||
const auto now = std::chrono::system_clock::now().time_since_epoch();
|
||||
return static_cast<uint64_t>(std::chrono::duration_cast<std::chrono::seconds>(now).count());
|
||||
return static_cast<uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(now).count());
|
||||
}()};
|
||||
return std::to_string(counter.fetch_add(1, std::memory_order_relaxed));
|
||||
}
|
||||
@@ -742,6 +743,10 @@ void OrcaPrinterAgent::forget_device_capabilities(const std::string& dev_id)
|
||||
|
||||
void OrcaPrinterAgent::deliver_to_sink(const std::string& dev_id, const std::string& payload, bool local)
|
||||
{
|
||||
// A files.* reply answers a pending request and is not printer state.
|
||||
if (local && try_consume_files_reply(dev_id, payload))
|
||||
return;
|
||||
|
||||
parse_ipcam_info(dev_id, payload);
|
||||
std::string merged_payload = merge_capabilities(dev_id, payload);
|
||||
|
||||
@@ -1218,7 +1223,7 @@ int OrcaPrinterAgent::connect_printer(const PrinterConnectionParams& params)
|
||||
m_camera_stream_mode = CameraStreamMode::none;
|
||||
m_camera_url.clear();
|
||||
m_current_connection = LAN;
|
||||
lan_mqtt_connection = std::make_unique<OrcaMqttConnection>();
|
||||
lan_mqtt_connection = std::make_shared<OrcaMqttConnection>();
|
||||
conn = lan_mqtt_connection.get();
|
||||
}
|
||||
BOOST_LOG_TRIVIAL(info) << "OrcaPrinterAgent: selected LAN printer dev_id=" << params.dev_id
|
||||
@@ -1279,7 +1284,7 @@ int OrcaPrinterAgent::disconnect_printer()
|
||||
{
|
||||
BOOST_LOG_TRIVIAL(info) << "Orca diagnostic: disconnect_printer requested";
|
||||
++m_lan_generation; // fence stale worker callbacks
|
||||
std::unique_ptr<OrcaMqttConnection> doomed;
|
||||
std::shared_ptr<OrcaMqttConnection> doomed;
|
||||
std::string prev_dev;
|
||||
CurrentConn previous_connection;
|
||||
CurrentConn current_connection;
|
||||
@@ -1759,35 +1764,48 @@ int OrcaPrinterAgent::start_send_gcode_to_sdcard(PrintParams params,
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
}
|
||||
|
||||
// Pure normalization of Moonraker's /server/files/list reply: the `result` array
|
||||
// of {path, modified, size}. `name` is the basename of `path`; malformed entries
|
||||
// (non-object, missing/empty path) are skipped. size/modified default to 0.
|
||||
std::vector<PrinterFileEntry> OrcaPrinterAgent::parse_file_list(const std::string& body)
|
||||
// Pure normalization of the files.list MQTT reply:
|
||||
// {"files":{...,"entries":[{name,path,is_dir,size,modified}]}}. Directory entries
|
||||
// and entries without a usable name are skipped; a missing path falls back to the
|
||||
// name and size/modified default to 0.
|
||||
std::vector<PrinterFileEntry> OrcaPrinterAgent::parse_files_list_reply(const std::string& payload)
|
||||
{
|
||||
std::vector<PrinterFileEntry> files;
|
||||
|
||||
const nlohmann::json envelope = nlohmann::json::parse(body, nullptr, false);
|
||||
const nlohmann::json envelope = nlohmann::json::parse(payload, nullptr, false);
|
||||
if (envelope.is_discarded() || !envelope.is_object())
|
||||
return files;
|
||||
|
||||
const auto result_it = envelope.find("result");
|
||||
if (result_it == envelope.end() || !result_it->is_array())
|
||||
const auto files_it = envelope.find("files");
|
||||
if (files_it == envelope.end() || !files_it->is_object())
|
||||
return files;
|
||||
|
||||
for (const auto& item : *result_it) {
|
||||
const auto entries_it = files_it->find("entries");
|
||||
if (entries_it == files_it->end() || !entries_it->is_array())
|
||||
return files;
|
||||
|
||||
for (const auto& item : *entries_it) {
|
||||
if (!item.is_object())
|
||||
continue;
|
||||
|
||||
const auto path_it = item.find("path");
|
||||
if (path_it == item.end() || !path_it->is_string())
|
||||
// Type-check is_dir rather than value(): a non-boolean would throw.
|
||||
const auto is_dir_it = item.find("is_dir");
|
||||
if (is_dir_it != item.end() && is_dir_it->is_boolean() && is_dir_it->get<bool>())
|
||||
continue;
|
||||
const std::string path = path_it->get<std::string>();
|
||||
if (path.empty())
|
||||
|
||||
const auto name_it = item.find("name");
|
||||
if (name_it == item.end() || !name_it->is_string())
|
||||
continue;
|
||||
const std::string name = name_it->get<std::string>();
|
||||
if (name.empty())
|
||||
continue;
|
||||
|
||||
PrinterFileEntry entry;
|
||||
entry.path = path;
|
||||
entry.name = fs::path(path).filename().string();
|
||||
entry.name = name;
|
||||
|
||||
const auto path_it = item.find("path");
|
||||
entry.path = (path_it != item.end() && path_it->is_string() && !path_it->get<std::string>().empty())
|
||||
? path_it->get<std::string>()
|
||||
: name;
|
||||
|
||||
const auto size_it = item.find("size");
|
||||
if (size_it != item.end() && size_it->is_number())
|
||||
@@ -1802,76 +1820,150 @@ std::vector<PrinterFileEntry> OrcaPrinterAgent::parse_file_list(const std::strin
|
||||
return files;
|
||||
}
|
||||
|
||||
int OrcaPrinterAgent::list_printer_files(const std::string& dev_id, PrinterFileListFn callback)
|
||||
// Dispatch one files.* command over the LAN MQTT connection and route its reply to
|
||||
// `complete`, matched by sequence_id. The UI-thread marshalling and the timing-out
|
||||
// watchdog capture only values and the heap-owned registry, never `this`.
|
||||
int OrcaPrinterAgent::send_files_request(const std::string& dev_id, const std::string& command, nlohmann::json fields,
|
||||
std::function<void(int result, const std::string& payload)> complete)
|
||||
{
|
||||
std::string origin;
|
||||
bool use_ssl = false;
|
||||
std::string ca_file;
|
||||
QueueOnMainFn queue;
|
||||
bool live = false;
|
||||
if (!complete)
|
||||
return BAMBU_NETWORK_ERR_INVALID_HANDLE;
|
||||
|
||||
std::shared_ptr<OrcaMqttConnection> conn;
|
||||
QueueOnMainFn queue;
|
||||
CurrentConn transport = NONE;
|
||||
std::string lan_dev;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(state_mutex);
|
||||
live = m_current_connection == LAN && m_lan_dev_id == dev_id;
|
||||
if (live) {
|
||||
origin = http_origin_from_lan_ws(m_lan_url);
|
||||
use_ssl = m_lan_use_ssl;
|
||||
ca_file = m_lan_ca_file;
|
||||
queue = queue_on_main_fn;
|
||||
}
|
||||
transport = m_current_connection;
|
||||
lan_dev = m_lan_dev_id;
|
||||
if (m_current_connection == LAN && m_lan_dev_id == dev_id)
|
||||
conn = lan_mqtt_connection; // keep alive across the send, racing disconnects included
|
||||
queue = queue_on_main_fn;
|
||||
}
|
||||
if (!live) {
|
||||
if (callback)
|
||||
callback(ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED, {});
|
||||
if (!conn) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files." << command << " not sent: no live LAN session for dev_id=" << dev_id
|
||||
<< " transport=" << connection_type_name(transport) << " m_lan_dev_id=" << lan_dev;
|
||||
complete(ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED, {});
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
if (origin.empty()) {
|
||||
if (callback)
|
||||
callback(BAMBU_NETWORK_ERR_INVALID_HANDLE, {});
|
||||
return BAMBU_NETWORK_ERR_INVALID_HANDLE;
|
||||
}
|
||||
|
||||
// perform_sync blocks, so the request runs off the UI thread. The worker captures
|
||||
// only values (never `this`). A trusted LAN facade needs no API key.
|
||||
std::thread([dev_id, origin, use_ssl, ca_file, queue, callback = std::move(callback)]() mutable {
|
||||
std::string body;
|
||||
int result = BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
fields["command"] = command;
|
||||
fields["sequence_id"] = next_command_sequence_id();
|
||||
const std::string sequence_id = fields["sequence_id"].get<std::string>();
|
||||
const std::string payload = nlohmann::json{{"files", std::move(fields)}}.dump();
|
||||
|
||||
auto http = Http::get(origin + "/server/files/list?root=gcodes");
|
||||
http.tls_verify(use_ssl);
|
||||
if (!ca_file.empty())
|
||||
http.ca_file(ca_file);
|
||||
http.timeout_connect(5)
|
||||
.timeout_max(15)
|
||||
.on_complete([&](std::string b, unsigned status) {
|
||||
if (status == 200) {
|
||||
body = std::move(b);
|
||||
result = BAMBU_NETWORK_SUCCESS;
|
||||
}
|
||||
})
|
||||
.on_error([&](std::string, std::string err, unsigned status) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "OrcaPrinterAgent: file list request failed status=" << status << " err=" << err;
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
std::vector<PrinterFileEntry> files;
|
||||
if (result == BAMBU_NETWORK_SUCCESS) {
|
||||
files = parse_file_list(body);
|
||||
// Empty is a valid listing; an unparseable body is not.
|
||||
if (nlohmann::json::parse(body, nullptr, false).is_discarded())
|
||||
result = BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
}
|
||||
|
||||
if (!callback)
|
||||
auto registry = m_files_replies;
|
||||
auto fired = std::make_shared<std::atomic<bool>>(false);
|
||||
auto finish = [queue, complete = std::move(complete), fired](int result, const std::string& reply) {
|
||||
if (fired->exchange(true))
|
||||
return;
|
||||
if (queue)
|
||||
queue([callback, result, files = std::move(files)]() mutable { callback(result, std::move(files)); });
|
||||
queue([complete, result, reply]() { complete(result, reply); });
|
||||
else
|
||||
callback(result, std::move(files));
|
||||
complete(result, reply);
|
||||
};
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(registry->mutex);
|
||||
registry->entries[sequence_id] = FilesReply{dev_id, finish, fired};
|
||||
}
|
||||
|
||||
if (!conn->send_request(dev_id, payload)) {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(registry->mutex);
|
||||
registry->entries.erase(sequence_id);
|
||||
}
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files." << command << " send_request failed dev_id=" << dev_id
|
||||
<< " sequence_id=" << sequence_id << " mqtt_connected=" << conn->is_connected();
|
||||
finish(BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED, {});
|
||||
return BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
}
|
||||
|
||||
// Watchdog: the connector answers promptly on a healthy link; fail the request
|
||||
// if nothing arrives so the UI does not wait on a dropped report.
|
||||
std::thread([registry, sequence_id, command, fired, finish]() {
|
||||
for (int waited = 0; waited < 100 && !fired->load(); ++waited)
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
||||
if (fired->load())
|
||||
return;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(registry->mutex);
|
||||
if (registry->entries.erase(sequence_id) == 0)
|
||||
return; // consumed concurrently
|
||||
}
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files." << command << " timed out waiting for reply sequence_id=" << sequence_id;
|
||||
finish(BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED, {});
|
||||
}).detach();
|
||||
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
}
|
||||
|
||||
bool OrcaPrinterAgent::try_consume_files_reply(const std::string& dev_id, const std::string& payload)
|
||||
{
|
||||
if (payload.find("\"files\"") == std::string::npos)
|
||||
return false;
|
||||
|
||||
const nlohmann::json envelope = nlohmann::json::parse(payload, nullptr, false);
|
||||
if (envelope.is_discarded() || !envelope.is_object())
|
||||
return false;
|
||||
|
||||
const auto files_it = envelope.find("files");
|
||||
if (files_it == envelope.end() || !files_it->is_object())
|
||||
return false;
|
||||
|
||||
const auto sequence_it = files_it->find("sequence_id");
|
||||
if (sequence_it == files_it->end() || !sequence_it->is_string()) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files frame without string sequence_id dev_id=" << dev_id;
|
||||
return false;
|
||||
}
|
||||
const std::string sequence_id = sequence_it->get<std::string>();
|
||||
|
||||
std::function<void(int, const std::string&)> complete;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(m_files_replies->mutex);
|
||||
const auto it = m_files_replies->entries.find(sequence_id);
|
||||
if (it == m_files_replies->entries.end()) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files reply sequence_id=" << sequence_id
|
||||
<< " has no pending request dev_id=" << dev_id;
|
||||
return false;
|
||||
}
|
||||
if (it->second.device_id != dev_id) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files reply sequence_id=" << sequence_id << " device mismatch: pending="
|
||||
<< it->second.device_id << " reply=" << dev_id;
|
||||
return false;
|
||||
}
|
||||
complete = it->second.complete;
|
||||
m_files_replies->entries.erase(it);
|
||||
}
|
||||
|
||||
// A non-string result (e.g. null) is a failure, never a throw.
|
||||
const auto result_it = files_it->find("result");
|
||||
const bool ok = result_it != files_it->end() && result_it->is_string() && result_it->get<std::string>() == "success";
|
||||
if (!ok) {
|
||||
const auto reason_it = files_it->find("reason");
|
||||
BOOST_LOG_TRIVIAL(warning) << "Orca diagnostic: files reply result=fail sequence_id=" << sequence_id
|
||||
<< " reason=" << (reason_it != files_it->end() && reason_it->is_string() ? reason_it->get<std::string>() : std::string());
|
||||
}
|
||||
complete(ok ? BAMBU_NETWORK_SUCCESS : BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED, ok ? payload : std::string());
|
||||
return true;
|
||||
}
|
||||
|
||||
// List the printer's G-code files over the LAN MQTT files.list command — the same
|
||||
// transport OrcaCloud uses, so the list works whether or not OrcaSonar's Moonraker
|
||||
// façade is enabled.
|
||||
int OrcaPrinterAgent::list_printer_files(const std::string& dev_id, PrinterFileListFn callback)
|
||||
{
|
||||
return send_files_request(
|
||||
dev_id, "list", nlohmann::json{{"root", "gcodes"}, {"path", ""}},
|
||||
[callback = std::move(callback)](int result, const std::string& payload) mutable {
|
||||
if (!callback)
|
||||
return;
|
||||
callback(result, result == BAMBU_NETWORK_SUCCESS ? parse_files_list_reply(payload)
|
||||
: std::vector<PrinterFileEntry>{});
|
||||
});
|
||||
}
|
||||
|
||||
// Pure pick of the widest thumbnail path from Moonraker's /server/files/thumbnails
|
||||
// reply. The `result` array is ordered smallest-first, so the largest `width` is
|
||||
// chosen. The key was renamed across Moonraker versions; both spellings are accepted.
|
||||
@@ -1905,31 +1997,32 @@ std::string OrcaPrinterAgent::parse_thumbnail_path(const std::string& body)
|
||||
return path;
|
||||
}
|
||||
|
||||
// Pure normalization of Moonraker's /server/files/metadata reply: the `result`
|
||||
// object's estimated_time (seconds), filament_total (mm) and filament_weight_total
|
||||
// (grams). A malformed reply or any non-numeric field defaults to 0.
|
||||
PrinterFileMetadata OrcaPrinterAgent::parse_file_metadata(const std::string& body)
|
||||
// Pure normalization of the files.metadata MQTT reply: the `files` object's
|
||||
// estimated_time (seconds), filament_total (mm) and filament_weight_total (grams;
|
||||
// OrcaSonar currently reports only filament_total). A malformed reply or any
|
||||
// non-numeric field defaults to 0.
|
||||
PrinterFileMetadata OrcaPrinterAgent::parse_files_metadata_reply(const std::string& payload)
|
||||
{
|
||||
PrinterFileMetadata meta;
|
||||
|
||||
const nlohmann::json envelope = nlohmann::json::parse(body, nullptr, false);
|
||||
const nlohmann::json envelope = nlohmann::json::parse(payload, nullptr, false);
|
||||
if (envelope.is_discarded() || !envelope.is_object())
|
||||
return meta;
|
||||
|
||||
const auto result_it = envelope.find("result");
|
||||
if (result_it == envelope.end() || !result_it->is_object())
|
||||
const auto files_it = envelope.find("files");
|
||||
if (files_it == envelope.end() || !files_it->is_object())
|
||||
return meta;
|
||||
|
||||
const auto time_it = result_it->find("estimated_time");
|
||||
if (time_it != result_it->end() && time_it->is_number())
|
||||
const auto time_it = files_it->find("estimated_time");
|
||||
if (time_it != files_it->end() && time_it->is_number())
|
||||
meta.estimated_time = static_cast<int>(time_it->get<double>());
|
||||
|
||||
const auto total_it = result_it->find("filament_total");
|
||||
if (total_it != result_it->end() && total_it->is_number())
|
||||
const auto total_it = files_it->find("filament_total");
|
||||
if (total_it != files_it->end() && total_it->is_number())
|
||||
meta.filament_total = total_it->get<double>();
|
||||
|
||||
const auto weight_it = result_it->find("filament_weight_total");
|
||||
if (weight_it != result_it->end() && weight_it->is_number())
|
||||
const auto weight_it = files_it->find("filament_weight_total");
|
||||
if (weight_it != files_it->end() && weight_it->is_number())
|
||||
meta.filament_weight = weight_it->get<double>();
|
||||
|
||||
return meta;
|
||||
@@ -2050,131 +2143,27 @@ int OrcaPrinterAgent::get_printer_file_thumbnail(const std::string& dev_id, cons
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
}
|
||||
|
||||
// Delete one G-code file via Moonraker's HTTP DELETE endpoint. Mirrors
|
||||
// list_printer_files for the connection snapshot and off-thread marshalling; the
|
||||
// worker captures no `this`.
|
||||
// Delete one G-code file over the LAN MQTT files.delete command.
|
||||
int OrcaPrinterAgent::delete_printer_file(const std::string& dev_id, const std::string& path, PrinterFileDeleteFn callback)
|
||||
{
|
||||
std::string origin;
|
||||
bool use_ssl = false;
|
||||
std::string ca_file;
|
||||
QueueOnMainFn queue;
|
||||
bool live = false;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(state_mutex);
|
||||
live = m_current_connection == LAN && m_lan_dev_id == dev_id;
|
||||
if (live) {
|
||||
origin = http_origin_from_lan_ws(m_lan_url);
|
||||
use_ssl = m_lan_use_ssl;
|
||||
ca_file = m_lan_ca_file;
|
||||
queue = queue_on_main_fn;
|
||||
}
|
||||
}
|
||||
if (!live) {
|
||||
if (callback)
|
||||
callback(ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED);
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
if (origin.empty()) {
|
||||
if (callback)
|
||||
callback(BAMBU_NETWORK_ERR_INVALID_HANDLE);
|
||||
return BAMBU_NETWORK_ERR_INVALID_HANDLE;
|
||||
}
|
||||
|
||||
std::thread([path, origin, use_ssl, ca_file, queue, callback = std::move(callback)]() mutable {
|
||||
int result = BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
|
||||
auto http = Http::del(origin + "/server/files/gcodes/" + encode_file_path(path));
|
||||
http.tls_verify(use_ssl);
|
||||
if (!ca_file.empty())
|
||||
http.ca_file(ca_file);
|
||||
http.timeout_connect(5)
|
||||
.timeout_max(15)
|
||||
.on_complete([&](std::string, unsigned status) {
|
||||
if (status == 200)
|
||||
result = BAMBU_NETWORK_SUCCESS;
|
||||
})
|
||||
.on_error([&](std::string, std::string err, unsigned status) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "OrcaPrinterAgent: file delete request failed status=" << status << " err=" << err;
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
if (!callback)
|
||||
return;
|
||||
if (queue)
|
||||
queue([callback, result]() mutable { callback(result); });
|
||||
else
|
||||
callback(result);
|
||||
}).detach();
|
||||
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
return send_files_request(dev_id, "delete", nlohmann::json{{"root", "gcodes"}, {"path", path}},
|
||||
[callback = std::move(callback)](int result, const std::string&) mutable {
|
||||
if (callback)
|
||||
callback(result);
|
||||
});
|
||||
}
|
||||
|
||||
// Fetch one file's Moonraker metadata (print time and filament usage). Mirrors
|
||||
// list_printer_files for the connection snapshot and off-thread marshalling; the
|
||||
// worker captures no `this`.
|
||||
// Fetch one file's metadata (print time and filament usage) over the LAN MQTT
|
||||
// files.metadata command.
|
||||
int OrcaPrinterAgent::get_printer_file_metadata(const std::string& dev_id, const std::string& path, PrinterFileMetadataFn callback)
|
||||
{
|
||||
std::string origin;
|
||||
bool use_ssl = false;
|
||||
std::string ca_file;
|
||||
QueueOnMainFn queue;
|
||||
bool live = false;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(state_mutex);
|
||||
live = m_current_connection == LAN && m_lan_dev_id == dev_id;
|
||||
if (live) {
|
||||
origin = http_origin_from_lan_ws(m_lan_url);
|
||||
use_ssl = m_lan_use_ssl;
|
||||
ca_file = m_lan_ca_file;
|
||||
queue = queue_on_main_fn;
|
||||
}
|
||||
}
|
||||
if (!live) {
|
||||
if (callback)
|
||||
callback(ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED, {});
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
if (origin.empty()) {
|
||||
if (callback)
|
||||
callback(BAMBU_NETWORK_ERR_INVALID_HANDLE, {});
|
||||
return BAMBU_NETWORK_ERR_INVALID_HANDLE;
|
||||
}
|
||||
|
||||
std::thread([path, origin, use_ssl, ca_file, queue, callback = std::move(callback)]() mutable {
|
||||
PrinterFileMetadata meta;
|
||||
int result = BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
|
||||
auto http = Http::get(origin + "/server/files/metadata?filename=" + Http::url_encode(path));
|
||||
http.tls_verify(use_ssl);
|
||||
if (!ca_file.empty())
|
||||
http.ca_file(ca_file);
|
||||
http.timeout_connect(5)
|
||||
.timeout_max(15)
|
||||
.on_complete([&](std::string b, unsigned status) {
|
||||
if (status == 200) {
|
||||
if (nlohmann::json::parse(b, nullptr, false).is_discarded())
|
||||
result = BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED;
|
||||
else {
|
||||
meta = parse_file_metadata(b);
|
||||
result = BAMBU_NETWORK_SUCCESS;
|
||||
}
|
||||
}
|
||||
})
|
||||
.on_error([&](std::string, std::string err, unsigned status) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "OrcaPrinterAgent: file metadata request failed status=" << status << " err=" << err;
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
if (!callback)
|
||||
return;
|
||||
if (queue)
|
||||
queue([callback, result, meta]() mutable { callback(result, meta); });
|
||||
else
|
||||
callback(result, meta);
|
||||
}).detach();
|
||||
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
return send_files_request(
|
||||
dev_id, "metadata", nlohmann::json{{"root", "gcodes"}, {"path", path}},
|
||||
[callback = std::move(callback)](int result, const std::string& payload) mutable {
|
||||
if (!callback)
|
||||
return;
|
||||
callback(result, result == BAMBU_NETWORK_SUCCESS ? parse_files_metadata_reply(payload) : PrinterFileMetadata{});
|
||||
});
|
||||
}
|
||||
|
||||
// Upload the sliced G-code, then start it: the LAN "print now" path.
|
||||
@@ -2279,7 +2268,7 @@ int OrcaPrinterAgent::start_sdcard_print(PrintParams params, OnUpdateStatusFn up
|
||||
BOOST_LOG_TRIVIAL(info) << "OrcaPrinterAgent: start_sdcard_print emitting filament_mapping entries=" << filament_mapping.size();
|
||||
}
|
||||
|
||||
nlohmann::json j = build_gcode_file_payload(next_gcode_file_sequence_id(), target, filament_mapping);
|
||||
nlohmann::json j = build_gcode_file_payload(next_command_sequence_id(), target, filament_mapping);
|
||||
|
||||
if (update_fn)
|
||||
update_fn(PrintingStageSending, 0, "Starting print...");
|
||||
|
||||
@@ -13,6 +13,7 @@
|
||||
#include <mutex>
|
||||
#include <memory>
|
||||
#include <thread>
|
||||
#include <unordered_map>
|
||||
#include <vector>
|
||||
|
||||
namespace Slic3r { class ICloudServiceAgent; }
|
||||
@@ -170,17 +171,31 @@ protected:
|
||||
const std::string& target,
|
||||
const nlohmann::json& filament_mapping);
|
||||
|
||||
// Pure JSON -> entries normalization for OrcaSonar's /server/files/list reply.
|
||||
// protected static so the test Probe reaches it.
|
||||
static std::vector<PrinterFileEntry> parse_file_list(const std::string& body);
|
||||
// Pure JSON -> entries normalization for the files.list MQTT reply
|
||||
// ({"files":{...,"entries":[{name,path,is_dir,size,modified}]}}). Directory
|
||||
// entries and entries without a usable name are skipped; a missing path falls
|
||||
// back to the name and size/modified default to 0. protected static for the Probe.
|
||||
static std::vector<PrinterFileEntry> parse_files_list_reply(const std::string& payload);
|
||||
|
||||
// Pick the widest thumbnail path from OrcaSonar's /server/files/thumbnails
|
||||
// reply (the array is smallest-first). Empty when none carry a path.
|
||||
static std::string parse_thumbnail_path(const std::string& body);
|
||||
|
||||
// Pure JSON -> metadata normalization for OrcaSonar's /server/files/metadata
|
||||
// reply. Missing or malformed fields default to 0. protected static for the Probe.
|
||||
static PrinterFileMetadata parse_file_metadata(const std::string& body);
|
||||
// Pure JSON -> metadata normalization for the files.metadata MQTT reply. Missing
|
||||
// or malformed fields default to 0. protected static for the Probe.
|
||||
static PrinterFileMetadata parse_files_metadata_reply(const std::string& payload);
|
||||
|
||||
// Dispatch one files.* command (list/metadata/delete) over the LAN MQTT
|
||||
// connection and route its reply, matched by sequence_id, into `complete`.
|
||||
// `complete` runs once on the UI thread (or the calling thread when no queue is
|
||||
// set) with the raw reply payload on success, or an empty string on
|
||||
// failure/timeout. Returns whether the request was dispatched.
|
||||
int send_files_request(const std::string& dev_id, const std::string& command, nlohmann::json fields,
|
||||
std::function<void(int result, const std::string& payload)> complete);
|
||||
|
||||
// Consume an inbound report that answers a pending files.* request. Returns true
|
||||
// when it matched, so the caller does not forward it to the machine-state sink.
|
||||
bool try_consume_files_reply(const std::string& dev_id, const std::string& payload);
|
||||
|
||||
// Percent-encode each '/'-separated segment for a Moonraker URL while keeping
|
||||
// the separators intact. protected static for the test Probe.
|
||||
@@ -218,7 +233,7 @@ private:
|
||||
CurrentConn m_current_connection = NONE;
|
||||
|
||||
std::shared_ptr<ICloudServiceAgent> m_cloud_agent;
|
||||
std::unique_ptr<OrcaMqttConnection> lan_mqtt_connection;
|
||||
std::shared_ptr<OrcaMqttConnection> lan_mqtt_connection;
|
||||
|
||||
// Two independent epochs: a cloud (de)selection must not fence the live LAN
|
||||
// feed, and vice versa. Each transport's connect thread and inbound handler
|
||||
@@ -241,6 +256,20 @@ private:
|
||||
CameraStreamMode m_camera_stream_mode = CameraStreamMode::none; // guarded by state_mutex
|
||||
std::string m_camera_url; // guarded by state_mutex
|
||||
|
||||
// Pending files.* MQTT replies, keyed by sequence_id. The registry is
|
||||
// heap-owned so a timeout watchdog can outlive the agent without touching
|
||||
// `this`; each entry carries its own one-shot guard.
|
||||
struct FilesReply {
|
||||
std::string device_id;
|
||||
std::function<void(int, const std::string&)> complete;
|
||||
std::shared_ptr<std::atomic<bool>> fired = std::make_shared<std::atomic<bool>>(false);
|
||||
};
|
||||
struct FilesReplyRegistry {
|
||||
std::mutex mutex;
|
||||
std::unordered_map<std::string, FilesReply> entries;
|
||||
};
|
||||
std::shared_ptr<FilesReplyRegistry> m_files_replies = std::make_shared<FilesReplyRegistry>();
|
||||
|
||||
OrcaCloudServiceAgent* get_orca_cloud_agent();
|
||||
|
||||
OrcaMqttConnection* get_appropriate_mqtt_connection(bool is_lan = true);
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
@@ -31,10 +31,11 @@ struct Probe : OrcaPrinterAgent {
|
||||
using OrcaPrinterAgent::build_filament_mapping;
|
||||
using OrcaPrinterAgent::build_gcode_file_payload;
|
||||
using OrcaPrinterAgent::prepare_outgoing_request;
|
||||
using OrcaPrinterAgent::parse_file_list;
|
||||
using OrcaPrinterAgent::parse_files_list_reply;
|
||||
using OrcaPrinterAgent::parse_thumbnail_path;
|
||||
using OrcaPrinterAgent::encode_file_path;
|
||||
using OrcaPrinterAgent::parse_file_metadata;
|
||||
using OrcaPrinterAgent::parse_files_metadata_reply;
|
||||
using OrcaPrinterAgent::try_consume_files_reply;
|
||||
};
|
||||
}
|
||||
|
||||
@@ -182,12 +183,17 @@ TEST_CASE("OrcaPrinterAgent::make_lan_client_id is stable and prefixed", "[OrcaP
|
||||
CHECK(a.rfind("orcaslicer-lan-dev-1-", 0) == 0);
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_list normalizes Moonraker entries", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files = Probe::parse_file_list(R"({
|
||||
"result": [
|
||||
{"path": "sub/foo.gcode", "modified": 1700000000.75, "size": 1234, "permissions": "rw"},
|
||||
{"path": "bar.gcode", "modified": 42, "size": 7}
|
||||
]
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_list_reply normalizes the MQTT entry fields", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files = Probe::parse_files_list_reply(R"({
|
||||
"files": {
|
||||
"command": "list",
|
||||
"sequence_id": "70001",
|
||||
"result": "success",
|
||||
"entries": [
|
||||
{"name": "foo.gcode", "path": "sub/foo.gcode", "is_dir": false, "size": 1234, "modified": 1700000000.75},
|
||||
{"name": "bar.gcode", "path": "bar.gcode", "is_dir": false, "size": 7, "modified": 42}
|
||||
]
|
||||
}
|
||||
})");
|
||||
REQUIRE(files.size() == 2);
|
||||
CHECK(files[0].path == "sub/foo.gcode");
|
||||
@@ -200,31 +206,54 @@ TEST_CASE("OrcaPrinterAgent::parse_file_list normalizes Moonraker entries", "[Or
|
||||
CHECK(files[1].modified == 42);
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_list handles an empty result", "[OrcaPrinterAgent]") {
|
||||
CHECK(Probe::parse_file_list(R"({"result": []})").empty());
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_list rejects malformed JSON", "[OrcaPrinterAgent]") {
|
||||
CHECK(Probe::parse_file_list("not json").empty());
|
||||
CHECK(Probe::parse_file_list(R"({"result": "nope"})").empty());
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_list skips entries missing fields", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files = Probe::parse_file_list(R"({
|
||||
"result": [
|
||||
{"size": 5},
|
||||
{"path": ""},
|
||||
{"path": "kept.gcode"},
|
||||
"not-an-object"
|
||||
]
|
||||
})");
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_list_reply falls back to the name when path is absent", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files =
|
||||
Probe::parse_files_list_reply(R"({"files": {"result": "success", "entries": [{"name": "solo.gcode"}]}})");
|
||||
REQUIRE(files.size() == 1);
|
||||
CHECK(files[0].path == "kept.gcode");
|
||||
CHECK(files[0].name == "kept.gcode");
|
||||
CHECK(files[0].path == "solo.gcode");
|
||||
CHECK(files[0].name == "solo.gcode");
|
||||
CHECK(files[0].size == 0);
|
||||
CHECK(files[0].modified == 0);
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_list_reply skips directories and unusable entries", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files = Probe::parse_files_list_reply(R"({
|
||||
"files": {
|
||||
"result": "success",
|
||||
"entries": [
|
||||
{"name": "sub", "path": "sub", "is_dir": true},
|
||||
{"path": "no-name.gcode"},
|
||||
{"name": ""},
|
||||
"not-an-object",
|
||||
{"name": "kept.gcode", "path": "kept.gcode"}
|
||||
]
|
||||
}
|
||||
})");
|
||||
REQUIRE(files.size() == 1);
|
||||
CHECK(files[0].name == "kept.gcode");
|
||||
CHECK(files[0].path == "kept.gcode");
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_list_reply rejects a malformed or empty reply", "[OrcaPrinterAgent]") {
|
||||
CHECK(Probe::parse_files_list_reply("not json").empty());
|
||||
CHECK(Probe::parse_files_list_reply(R"({"result": []})").empty());
|
||||
CHECK(Probe::parse_files_list_reply(R"({"files": {"entries": "nope"}})").empty());
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_list_reply tolerates a non-boolean is_dir", "[OrcaPrinterAgent]") {
|
||||
const std::vector<Slic3r::PrinterFileEntry> files =
|
||||
Probe::parse_files_list_reply(R"({"files": {"result": "success", "entries": [{"name": "x.gcode", "is_dir": "yes"}]}})");
|
||||
REQUIRE(files.size() == 1);
|
||||
CHECK(files[0].name == "x.gcode");
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::try_consume_files_reply ignores reports it did not request", "[OrcaPrinterAgent]") {
|
||||
Probe agent("/tmp");
|
||||
CHECK_FALSE(agent.try_consume_files_reply("dev-1", R"({"print": {"command": "push_status"}})"));
|
||||
CHECK_FALSE(agent.try_consume_files_reply("dev-1", R"({"files": {"command": "list"}})")); // no sequence_id
|
||||
CHECK_FALSE(agent.try_consume_files_reply("dev-1", R"({"files": {"command": "list", "sequence_id": "999"}})")); // not pending
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_thumbnail_path picks the largest width", "[OrcaPrinterAgent]") {
|
||||
const std::string path = Probe::parse_thumbnail_path(R"({
|
||||
"result": [
|
||||
@@ -263,17 +292,19 @@ TEST_CASE("OrcaPrinterAgent::encode_file_path preserves separators and encodes s
|
||||
CHECK(Probe::encode_file_path("design+part.gcode") == "design%2Bpart.gcode");
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_metadata parses the Moonraker fields", "[OrcaPrinterAgent]") {
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_metadata_reply parses the metadata fields", "[OrcaPrinterAgent]") {
|
||||
using Catch::Matchers::WithinAbs;
|
||||
const Slic3r::PrinterFileMetadata meta = Probe::parse_file_metadata(R"({
|
||||
"result": {
|
||||
"filename": "sub/foo.gcode",
|
||||
const Slic3r::PrinterFileMetadata meta = Probe::parse_files_metadata_reply(R"({
|
||||
"files": {
|
||||
"command": "metadata",
|
||||
"sequence_id": "70002",
|
||||
"result": "success",
|
||||
"path": "sub/foo.gcode",
|
||||
"size": 1234,
|
||||
"modified": 1700000000.5,
|
||||
"estimated_time": 3725,
|
||||
"filament_total": 10500.5,
|
||||
"filament_weight_total": 31.6,
|
||||
"thumbnails": []
|
||||
"filament_weight_total": 31.6
|
||||
}
|
||||
})");
|
||||
CHECK(meta.estimated_time == 3725);
|
||||
@@ -281,22 +312,22 @@ TEST_CASE("OrcaPrinterAgent::parse_file_metadata parses the Moonraker fields", "
|
||||
CHECK_THAT(meta.filament_weight, WithinAbs(31.6, 1e-9));
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_metadata defaults missing fields to zero", "[OrcaPrinterAgent]") {
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_metadata_reply defaults missing fields to zero", "[OrcaPrinterAgent]") {
|
||||
using Catch::Matchers::WithinAbs;
|
||||
const Slic3r::PrinterFileMetadata meta = Probe::parse_file_metadata(R"({"result": {"filename": "foo.gcode"}})");
|
||||
const Slic3r::PrinterFileMetadata meta = Probe::parse_files_metadata_reply(R"({"files": {"path": "foo.gcode"}})");
|
||||
CHECK(meta.estimated_time == 0);
|
||||
CHECK_THAT(meta.filament_total, WithinAbs(0.0, 1e-12));
|
||||
CHECK_THAT(meta.filament_weight, WithinAbs(0.0, 1e-12));
|
||||
}
|
||||
|
||||
TEST_CASE("OrcaPrinterAgent::parse_file_metadata yields zeros for a malformed reply", "[OrcaPrinterAgent]") {
|
||||
TEST_CASE("OrcaPrinterAgent::parse_files_metadata_reply yields zeros for a malformed reply", "[OrcaPrinterAgent]") {
|
||||
using Catch::Matchers::WithinAbs;
|
||||
const Slic3r::PrinterFileMetadata malformed = Probe::parse_file_metadata("not json");
|
||||
const Slic3r::PrinterFileMetadata malformed = Probe::parse_files_metadata_reply("not json");
|
||||
CHECK(malformed.estimated_time == 0);
|
||||
CHECK_THAT(malformed.filament_total, WithinAbs(0.0, 1e-12));
|
||||
CHECK_THAT(malformed.filament_weight, WithinAbs(0.0, 1e-12));
|
||||
|
||||
const Slic3r::PrinterFileMetadata wrong_type = Probe::parse_file_metadata(R"({"result": "nope"})");
|
||||
const Slic3r::PrinterFileMetadata wrong_type = Probe::parse_files_metadata_reply(R"({"files": "nope"})");
|
||||
CHECK(wrong_type.estimated_time == 0);
|
||||
CHECK_THAT(wrong_type.filament_total, WithinAbs(0.0, 1e-12));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user