Compare commits

...
3 changed files with 1259 additions and 151 deletions
File diff suppressed because it is too large Load Diff
+79 -3
View File
@@ -6,6 +6,7 @@
#include "OrcaCloudServiceAgent.hpp"
#include "OrcaMqttConnection.hpp"
#include <atomic>
#include <chrono>
#include <cstdint>
#include <functional>
#include "bambu_networking.hpp"
@@ -122,6 +123,12 @@ public:
void bump_lan_generation_for_test() { ++m_lan_generation; }
// Test-only: the same for the (independent) cloud selection epoch.
void bump_cloud_generation_for_test() { ++m_cloud_generation; }
// Test-only: shrink the OPCP file-request waits (probe and per-page).
void set_file_opcp_timeouts_for_test(int probe_ms, int request_ms)
{
m_file_probe_timeout_ms.store(probe_ms);
m_file_request_timeout_ms.store(request_ms);
}
FilamentSyncMode get_filament_sync_mode() const override { return FilamentSyncMode::subscription; }
@@ -142,8 +149,8 @@ protected:
// per-rule detail. Returns the payload unchanged when no rule applies.
std::string merge_capabilities(const std::string& dev_id, const std::string& payload);
void forget_device_capabilities(const std::string& dev_id);
int prepare_outgoing_request(const std::string& dev_id, const std::string& payload,
std::string& command, std::string& prepared) const;
static int prepare_outgoing_request(const std::string& dev_id, const std::string& payload,
std::string& command, std::string& prepared);
// Report the asynchronous LAN connection state using the same callback contract as
// the other printer agents. The transport result cannot be returned by
@@ -185,6 +192,65 @@ protected:
// Percent-encode each '/'-separated segment for a Moonraker URL while keeping
// the separators intact. protected static for the test Probe.
static std::string encode_file_path(const std::string& path);
// ========================================================================
// OPCP printer-file transport (Storage tab).
//
// Newer OrcaSonar builds expose files.* over the same MQTT channel the agent
// already owns; older ones only have the Moonraker HTTP façade. The manifest's
// files.list advertisement picks the transport, with a short OPCP probe and an
// HTTP fallback when the capability is not known yet.
// ========================================================================
// opcp: the manifest advertises files.list. http: the manifest is known and
// does not. probe: no manifest yet; try OPCP briefly then fall back to HTTP.
enum class FileTransport { opcp, http, probe };
static FileTransport preferred_file_transport(const std::string& dev_id);
// One files.list reply: the file entries the GUI shows, plus the subdirectory
// paths the walker still has to visit.
struct OpcpFileListPage {
std::vector<PrinterFileEntry> files;
std::vector<std::string> dirs;
bool truncated = false;
bool success = false;
};
static std::string build_files_list_request(const std::string& sequence_id, const std::string& path);
static std::string build_files_metadata_request(const std::string& sequence_id, const std::string& path);
static std::string build_files_delete_request(const std::string& sequence_id, const std::string& path);
static std::string build_files_thumbnail_request(const std::string& sequence_id, const std::string& path);
static OpcpFileListPage parse_opcp_file_list(const std::string& body);
static PrinterFileMetadata parse_opcp_file_metadata(const std::string& body);
static std::string parse_opcp_thumbnail(const std::string& body);
// True when a files.* reply is unusable, so the caller should try HTTP.
static bool file_reply_requires_http_fallback(const std::string& body);
// Consume an inbound files.* reply. Returns true for any files payload, which
// is then not forwarded to the GUI.
bool handle_files_reply(const std::string& dev_id, const std::string& payload);
// Wake every pending file request for a device (disconnect/unbind/forget).
void fail_pending_file_requests(const std::string& dev_id);
// Bounded BFS over OPCP files.list from the gcodes root. Returns false when the
// root listing itself failed, so the caller can fall back to HTTP.
static bool collect_opcp_file_entries(const std::function<int(const std::string&)>& send,
const std::string& dev_id,
std::chrono::milliseconds first_timeout,
std::chrono::milliseconds page_timeout,
std::vector<PrinterFileEntry>& out,
bool& truncated);
// HTTP fallbacks, split out of the public ops. Static so the detached workers
// that call them never capture `this`.
static void list_printer_files_http(const std::string& origin, bool use_ssl, const std::string& ca_file,
QueueOnMainFn queue, PrinterFileListFn callback);
static void get_printer_file_thumbnail_http(const std::string& path, const std::string& origin, bool use_ssl,
const std::string& ca_file, QueueOnMainFn queue, PrinterFileThumbnailFn callback);
static void delete_printer_file_http(const std::string& path, const std::string& origin, bool use_ssl,
const std::string& ca_file, QueueOnMainFn queue, PrinterFileDeleteFn callback);
static void get_printer_file_metadata_http(const std::string& path, const std::string& origin, bool use_ssl,
const std::string& ca_file, QueueOnMainFn queue, PrinterFileMetadataFn callback);
// Test hook: the ws:// URL connect_printer built for the current LAN session ("" if none).
std::string lan_connection_target() const;
// Shared post-connect sequence: SUBSCRIBE, then pushing.start, pushall,
@@ -218,7 +284,14 @@ private:
CurrentConn m_current_connection = NONE;
std::shared_ptr<ICloudServiceAgent> m_cloud_agent;
std::unique_ptr<OrcaMqttConnection> lan_mqtt_connection;
// shared_ptr so a detached file worker can hold the connection across a
// disconnect without a use-after-free; the worker never captures `this`.
std::shared_ptr<OrcaMqttConnection> lan_mqtt_connection;
// OPCP file-request waits, in milliseconds. Probe is short so an older build
// fails over to HTTP quickly; the per-page wait covers a real transfer.
std::atomic<int> m_file_probe_timeout_ms{3000};
std::atomic<int> m_file_request_timeout_ms{15000};
// 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
@@ -249,6 +322,9 @@ private:
// Route one command payload to device/<dev_id>/request on the LAN or the shared
// cloud connection. The uniform send path for both send_message* overrides.
int route_send(bool is_lan, const std::string& dev_id, const std::string& json_str);
// Core of route_send() on an already-resolved connection. Static so detached
// file workers (which must not capture `this`) send through the same path.
static int route_send_via(bool is_lan, OrcaMqttConnection* conn, const std::string& dev_id, const std::string& json_str);
// Callbacks
OnMsgArrivedFn on_ssdp_msg_fn;
@@ -9,9 +9,12 @@
#include "catch2/matchers/catch_matchers.hpp"
#include "orca_mqtt_mock_broker.hpp"
#include <atomic>
#include <chrono>
#include <functional>
#include <memory>
#include <mutex>
#include <set>
#include <slic3r/Utils/bambu_networking.hpp>
#include <string>
#include <thread>
@@ -35,7 +38,78 @@ struct Probe : OrcaPrinterAgent {
using OrcaPrinterAgent::parse_thumbnail_path;
using OrcaPrinterAgent::encode_file_path;
using OrcaPrinterAgent::parse_file_metadata;
using OrcaPrinterAgent::FileTransport;
using OrcaPrinterAgent::preferred_file_transport;
using OrcaPrinterAgent::build_files_list_request;
using OrcaPrinterAgent::build_files_metadata_request;
using OrcaPrinterAgent::build_files_delete_request;
using OrcaPrinterAgent::build_files_thumbnail_request;
using OrcaPrinterAgent::parse_opcp_file_list;
using OrcaPrinterAgent::parse_opcp_file_metadata;
using OrcaPrinterAgent::parse_opcp_thumbnail;
using OrcaPrinterAgent::file_reply_requires_http_fallback;
};
// Shared helpers for the OPCP file-transport tests below.
bool wait_for(const std::function<bool()>& predicate, int attempts = 400, int sleep_ms = 10)
{
for (int i = 0; i < attempts; ++i) {
if (predicate())
return true;
std::this_thread::sleep_for(std::chrono::milliseconds(sleep_ms));
}
return predicate();
}
// First recorded files.<command> request, optionally narrowed to one path.
std::string find_files_request(const orca_mqtt_test::MockBroker& broker, const std::string& command, const std::string& path)
{
for (const std::string& request : broker.received_requests()) {
if (request.find("\"files\"") == std::string::npos)
continue;
if (request.find("\"command\":\"" + command + "\"") == std::string::npos)
continue;
if (!path.empty() && request.find("\"path\":\"" + path + "\"") == std::string::npos)
continue;
return request;
}
return {};
}
std::string request_sequence_id(const std::string& request)
{
const nlohmann::json parsed = nlohmann::json::parse(request, nullptr, false);
if (parsed.is_discarded() || !parsed.is_object())
return {};
const auto files_it = parsed.find("files");
if (files_it == parsed.end() || !files_it->is_object())
return {};
return files_it->value("sequence_id", std::string());
}
std::string files_list_reply(const std::string& sequence_id, const std::string& entries_json, bool truncated = false)
{
return R"({"files":{"command":"list","sequence_id":")" + sequence_id + R"(","result":"success","errno":0,"entries":)" +
entries_json + R"(,"truncated":)" + (truncated ? "true" : "false") + R"(}})";
}
bool connect_and_subscribe(Probe& agent, const orca_mqtt_test::MockBroker& broker, const std::string& dev_id)
{
const auto endpoint = broker.host_port();
agent.connect_printer(
Slic3r::PrinterConnectionParams{dev_id, endpoint.first + ":" + endpoint.second, "orcasonar", "code", "", false, ""});
return wait_for([&] { return broker.connect_count() >= 1 && broker.is_subscribed("device/" + dev_id + "/report"); });
}
struct ListOutcome {
std::mutex mutex;
int result = 999;
std::vector<Slic3r::PrinterFileEntry> files;
};
const char* const kFilesListCapabilities =
R"({"info":{"command":"get_capabilities","supported_commands":["files.list"],"capabilities":{"protocol":{"supported_commands":[]}}}})";
}
TEST_CASE("OrcaPrinterAgent forwards a status payload to on_message_fn", "[OrcaPrinterAgent]") {
@@ -574,3 +648,331 @@ TEST_CASE("destroying an agent mid-connect does not hang or crash", "[OrcaPrinte
}
SUCCEED();
}
// ============================================================================
// OPCP printer-file transport
// ============================================================================
TEST_CASE("the file transport follows the advertised files.list capability", "[OrcaPrinterAgent]") {
Probe agent("/tmp");
// No manifest yet: probe OPCP briefly, then fall back to HTTP.
CHECK(Probe::preferred_file_transport("transport-unknown") == Probe::FileTransport::probe);
// Manifest advertises files.list: use OPCP.
agent.deliver_to_sink("transport-opcp", kFilesListCapabilities, true);
CHECK(Probe::preferred_file_transport("transport-opcp") == Probe::FileTransport::opcp);
// Manifest is known and lacks files.list: use HTTP.
agent.deliver_to_sink("transport-http",
R"({"info":{"command":"get_capabilities","supported_commands":["print.gcode_file"],"capabilities":{"protocol":{"supported_commands":[]}}}})",
true);
CHECK(Probe::preferred_file_transport("transport-http") == Probe::FileTransport::http);
}
TEST_CASE("the files.* request builders match the OPCP wire contract", "[OrcaPrinterAgent]") {
const nlohmann::json list = nlohmann::json::parse(Probe::build_files_list_request("7", ""));
CHECK(list["files"]["command"] == "list");
CHECK(list["files"]["sequence_id"] == "7");
CHECK(list["files"]["root"] == "gcodes");
CHECK(list["files"]["path"] == "");
const nlohmann::json metadata = nlohmann::json::parse(Probe::build_files_metadata_request("8", "sub/model.gcode"));
CHECK(metadata["files"]["command"] == "metadata");
CHECK(metadata["files"]["root"] == "gcodes");
CHECK(metadata["files"]["path"] == "sub/model.gcode");
const nlohmann::json del = nlohmann::json::parse(Probe::build_files_delete_request("9", "sub/model.gcode"));
CHECK(del["files"]["command"] == "delete");
CHECK(del["files"]["path"] == "sub/model.gcode");
const nlohmann::json thumb = nlohmann::json::parse(Probe::build_files_thumbnail_request("10", "sub/model.gcode"));
CHECK(thumb["files"]["command"] == "thumbnail");
CHECK(thumb["files"]["max_dimension"] == 0);
CHECK(thumb["files"]["path"] == "sub/model.gcode");
}
TEST_CASE("parse_opcp_file_list splits files from directories", "[OrcaPrinterAgent]") {
const auto page = Probe::parse_opcp_file_list(R"({
"files": {
"command": "list",
"sequence_id": "1",
"result": "success",
"errno": 0,
"root": "gcodes",
"path": "",
"entries": [
{"name": "a.gcode", "path": "a.gcode", "is_dir": false, "size": 10, "modified": 100.5},
{"name": "sub", "path": "sub", "is_dir": true, "size": 0, "modified": 0}
],
"truncated": false
}
})");
CHECK(page.success);
REQUIRE(page.files.size() == 1);
CHECK(page.files[0].path == "a.gcode");
CHECK(page.files[0].name == "a.gcode");
CHECK(page.files[0].size == 10);
CHECK(page.files[0].modified == 100);
REQUIRE(page.dirs.size() == 1);
CHECK(page.dirs[0] == "sub");
CHECK_FALSE(page.truncated);
}
TEST_CASE("parse_opcp_file_list reports truncation and failed replies", "[OrcaPrinterAgent]") {
const auto truncated = Probe::parse_opcp_file_list(R"({
"files": {"command": "list", "sequence_id": "1", "result": "success", "errno": 0, "entries": [], "truncated": true}
})");
CHECK(truncated.success);
CHECK(truncated.truncated);
const std::string failed_reply = R"({"files":{"command":"list","sequence_id":"1","result":"error","errno":-19,"entries":[]}})";
const auto failed = Probe::parse_opcp_file_list(failed_reply);
CHECK_FALSE(failed.success);
CHECK(Probe::file_reply_requires_http_fallback(failed_reply));
CHECK_FALSE(Probe::file_reply_requires_http_fallback(
R"({"files":{"command":"list","sequence_id":"1","result":"success","errno":0,"entries":[]}})"));
}
TEST_CASE("parse_opcp_file_metadata reads the OPCP fields", "[OrcaPrinterAgent]") {
using Catch::Matchers::WithinAbs;
const Slic3r::PrinterFileMetadata meta = Probe::parse_opcp_file_metadata(R"({
"files": {"command": "metadata", "result": "success", "errno": 0,
"estimated_time": 3725, "filament_total": 10500.5, "filament_weight_total": 31.6}
})");
CHECK(meta.estimated_time == 3725);
CHECK_THAT(meta.filament_total, WithinAbs(10500.5, 1e-9));
CHECK_THAT(meta.filament_weight, WithinAbs(31.6, 1e-9));
const Slic3r::PrinterFileMetadata empty = Probe::parse_opcp_file_metadata(R"({"files": {"command": "metadata"}})");
CHECK(empty.estimated_time == 0);
CHECK_THAT(empty.filament_total, WithinAbs(0.0, 1e-12));
CHECK_THAT(empty.filament_weight, WithinAbs(0.0, 1e-12));
}
TEST_CASE("parse_opcp_thumbnail decodes the base64 PNG payload", "[OrcaPrinterAgent]") {
// "aGVsbG8=" is base64 for "hello".
const std::string decoded = Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"width":300,"height":300,"size":5,"data":"aGVsbG8="}})");
CHECK(decoded == "hello");
// data absent means the file has no embedded thumbnail.
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"width":300,"height":300,"size":0}})").empty());
CHECK(Probe::parse_opcp_thumbnail("not json").empty());
}
TEST_CASE("parse_opcp_thumbnail rejects unpadded and malformed base64", "[OrcaPrinterAgent]") {
// Unpadded: Beast's decoded_size() requires a length multiple of four, so this
// must be rejected rather than decoded into an undersized buffer.
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"data":"aGVsbG8"}})").empty());
// Non-alphabet characters must not decode as a partial image (length 8 here,
// so this exercises the character check rather than the length check).
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"data":"aGVs*G8="}})").empty());
// Padding is only valid closing the final group.
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"data":"aG=s"}})").empty());
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"data":"===="}})").empty());
// A correctly padded payload still decodes.
CHECK(Probe::parse_opcp_thumbnail(
R"({"files":{"command":"thumbnail","result":"success","errno":0,"data":"aGVsbG8="}})") == "hello");
}
TEST_CASE("list_printer_files parses an OPCP files.list reply", "[OrcaPrinterAgent]") {
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-flat";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
auto outcome = std::make_shared<ListOutcome>();
agent.list_printer_files(dev_id, [outcome](int result, std::vector<Slic3r::PrinterFileEntry> files) {
std::lock_guard<std::mutex> lock(outcome->mutex);
outcome->result = result;
outcome->files = std::move(files);
});
std::string request;
REQUIRE(wait_for([&] { request = find_files_request(broker, "list", ""); return !request.empty(); }));
const std::string sequence_id = request_sequence_id(request);
REQUIRE_FALSE(sequence_id.empty());
broker.push_report(dev_id, files_list_reply(sequence_id,
R"([{"name":"a.gcode","path":"a.gcode","is_dir":false,"size":1234,"modified":1700000000.5}])"));
REQUIRE(wait_for([&] { std::lock_guard<std::mutex> lock(outcome->mutex); return outcome->result != 999; }));
{
std::lock_guard<std::mutex> lock(outcome->mutex);
CHECK(outcome->result == BAMBU_NETWORK_SUCCESS);
REQUIRE(outcome->files.size() == 1);
CHECK(outcome->files[0].name == "a.gcode");
CHECK(outcome->files[0].size == 1234);
CHECK(outcome->files[0].modified == 1700000000);
}
agent.disconnect_printer();
}
TEST_CASE("list_printer_files walks OPCP subdirectories", "[OrcaPrinterAgent]") {
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-recursive";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
auto outcome = std::make_shared<ListOutcome>();
agent.list_printer_files(dev_id, [outcome](int result, std::vector<Slic3r::PrinterFileEntry> files) {
std::lock_guard<std::mutex> lock(outcome->mutex);
outcome->result = result;
outcome->files = std::move(files);
});
// Reply to each files.list request the walker emits: root first, then "sub".
std::set<std::string> replied;
const bool finished = wait_for([&] {
{
std::lock_guard<std::mutex> lock(outcome->mutex);
if (outcome->result != 999)
return true;
}
for (const std::string& request : broker.received_requests()) {
if (request.find("\"command\":\"list\"") == std::string::npos)
continue;
const std::string sequence_id = request_sequence_id(request);
if (sequence_id.empty() || replied.count(sequence_id) != 0)
continue;
replied.insert(sequence_id);
const nlohmann::json parsed = nlohmann::json::parse(request, nullptr, false);
const std::string path = parsed.is_object() ? parsed["files"].value("path", std::string()) : std::string();
if (path.empty())
broker.push_report(dev_id, files_list_reply(sequence_id,
R"([{"name":"root.gcode","path":"root.gcode","is_dir":false,"size":11,"modified":100.5},)"
R"({"name":"sub","path":"sub","is_dir":true,"size":0,"modified":0}])"));
else if (path == "sub")
broker.push_report(dev_id, files_list_reply(sequence_id,
R"([{"name":"nested.gcode","path":"sub/nested.gcode","is_dir":false,"size":22,"modified":200.5}])"));
}
return false;
});
REQUIRE(finished);
{
std::lock_guard<std::mutex> lock(outcome->mutex);
CHECK(outcome->result == BAMBU_NETWORK_SUCCESS);
REQUIRE(outcome->files.size() == 2);
CHECK(outcome->files[0].path == "root.gcode");
CHECK(outcome->files[1].path == "sub/nested.gcode");
CHECK(outcome->files[1].size == 22);
}
agent.disconnect_printer();
}
TEST_CASE("an unanswered OPCP file listing times out and reports failure", "[OrcaPrinterAgent]") {
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-timeout";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
agent.set_file_opcp_timeouts_for_test(/*probe=*/150, /*request=*/150);
auto result = std::make_shared<std::atomic<int>>(999);
agent.list_printer_files(dev_id, [result](int r, std::vector<Slic3r::PrinterFileEntry>) { result->store(r); });
REQUIRE(wait_for([&] { return !find_files_request(broker, "list", "").empty(); }));
std::this_thread::sleep_for(std::chrono::milliseconds(300)); // let the request wait expire
agent.disconnect_printer(); // free the broker so the HTTP fallback fails fast
REQUIRE(wait_for([&] { return result->load() != 999; }, 500, 10));
CHECK(result->load() != BAMBU_NETWORK_SUCCESS);
}
TEST_CASE("disconnect_printer wakes a pending file request", "[OrcaPrinterAgent]") {
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-disconnect";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
// Far longer than the assertion window, so only fail_pending_file_requests
// can wake the worker.
agent.set_file_opcp_timeouts_for_test(/*probe=*/60000, /*request=*/60000);
auto result = std::make_shared<std::atomic<int>>(999);
agent.list_printer_files(dev_id, [result](int r, std::vector<Slic3r::PrinterFileEntry>) { result->store(r); });
REQUIRE(wait_for([&] { return !find_files_request(broker, "list", "").empty(); }));
REQUIRE(result->load() == 999); // still waiting before the disconnect
agent.disconnect_printer();
REQUIRE(wait_for([&] { return result->load() != 999; }, 200, 10)); // 2s, well under the 60s wait
CHECK(result->load() != BAMBU_NETWORK_SUCCESS);
}
TEST_CASE("an errno -19 files reply falls back to HTTP", "[OrcaPrinterAgent]") {
const std::string failed_reply = R"({"files":{"command":"list","sequence_id":"1","result":"error","errno":-19,"entries":[]}})";
CHECK(Probe::file_reply_requires_http_fallback(failed_reply));
CHECK_FALSE(Probe::parse_opcp_file_list(failed_reply).success);
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-errno";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
agent.set_file_opcp_timeouts_for_test(/*probe=*/500, /*request=*/500);
auto result = std::make_shared<std::atomic<int>>(999);
agent.list_printer_files(dev_id, [result](int r, std::vector<Slic3r::PrinterFileEntry>) { result->store(r); });
std::string request;
REQUIRE(wait_for([&] { request = find_files_request(broker, "list", ""); return !request.empty(); }));
const std::string sequence_id = request_sequence_id(request);
REQUIRE_FALSE(sequence_id.empty());
broker.push_report(dev_id,
R"({"files":{"command":"list","sequence_id":")" + sequence_id + R"(","result":"error","errno":-19,"entries":[]}})");
std::this_thread::sleep_for(std::chrono::milliseconds(150)); // let the worker consume the failed reply
agent.disconnect_printer(); // free the broker so the HTTP fallback fails fast
REQUIRE(wait_for([&] { return result->load() != 999; }, 500, 10));
CHECK(result->load() != BAMBU_NETWORK_SUCCESS);
}
TEST_CASE("files replies are consumed and never forwarded to the GUI", "[OrcaPrinterAgent]") {
orca_mqtt_test::MockBroker broker;
Probe agent("/tmp");
const std::string dev_id = "opcp-list-routing";
REQUIRE(connect_and_subscribe(agent, broker, dev_id));
agent.deliver_to_sink(dev_id, kFilesListCapabilities, true);
std::atomic<int> forwarded{0};
agent.set_on_local_message_fn([&forwarded](std::string, std::string) { forwarded.fetch_add(1); });
auto outcome = std::make_shared<ListOutcome>();
agent.list_printer_files(dev_id, [outcome](int result, std::vector<Slic3r::PrinterFileEntry> files) {
std::lock_guard<std::mutex> lock(outcome->mutex);
outcome->result = result;
outcome->files = std::move(files);
});
std::string request;
REQUIRE(wait_for([&] { request = find_files_request(broker, "list", ""); return !request.empty(); }));
const std::string sequence_id = request_sequence_id(request);
REQUIRE_FALSE(sequence_id.empty());
// The matching reply satisfies the pending request and is not forwarded.
agent.deliver_to_sink(dev_id,
files_list_reply(sequence_id, R"([{"name":"a.gcode","path":"a.gcode","is_dir":false,"size":1,"modified":1}])"), true);
REQUIRE(wait_for([&] { std::lock_guard<std::mutex> lock(outcome->mutex); return outcome->result != 999; }));
CHECK(forwarded.load() == 0);
// An unmatched files report is dropped too.
agent.deliver_to_sink(dev_id,
R"({"files":{"command":"list","sequence_id":"no-such-sequence","result":"success","errno":0,"entries":[]}})"), true);
CHECK(forwarded.load() == 0);
// An unrelated report still reaches the sink.
agent.deliver_to_sink(dev_id, R"({"print":{"command":"push_status"}})", true);
CHECK(forwarded.load() == 1);
agent.disconnect_printer();
}