mirror of
https://github.com/OrcaSlicer/OrcaSlicer.git
synced 2026-10-07 15:51:08 +00:00
fix: batch box mapping, de-dup fetches and reject unsupported starts
This commit is contained in:
@@ -589,13 +589,19 @@ int BBLPrinterAgent::start_local_print_with_record(PrintParams params, OnUpdateS
|
||||
|
||||
int BBLPrinterAgent::start_send_gcode_to_sdcard(PrintParams params, OnUpdateStatusFn update_fn, WasCancelledFn cancel_fn, OnWaitFn wait_fn)
|
||||
{
|
||||
// dispatch_start() moves out of `params`, so snapshot the diagnostic fields first;
|
||||
// logging them after the call would print empty strings.
|
||||
const bool try_emmc_print = params.try_emmc_print;
|
||||
const std::string dev_ip = params.dev_ip;
|
||||
const std::string dev_id = params.dev_id;
|
||||
|
||||
int result = dispatch_start<func_start_send_gcode_to_sdcard_legacy, func_start_send_gcode_to_sdcard_0203>(
|
||||
BBLNetworkPlugin::instance().get_start_send_gcode_to_sdcard(), params, update_fn, cancel_fn, wait_fn);
|
||||
if (result != 0) {
|
||||
BOOST_LOG_TRIVIAL(error) << "start_send_gcode_to_sdcard failed: result=" << result
|
||||
<< ", try_emmc_print=" << params.try_emmc_print
|
||||
<< ", try_emmc_print=" << try_emmc_print
|
||||
<< ", legacy_mode=" << BBLNetworkPlugin::instance().use_legacy_network()
|
||||
<< ", dev_ip=" << params.dev_ip << ", dev_id=" << params.dev_id;
|
||||
<< ", dev_ip=" << dev_ip << ", dev_id=" << dev_id;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -567,11 +567,13 @@ int MoonrakerPrinterAgent::set_user_selected_machine(std::string dev_id)
|
||||
|
||||
int MoonrakerPrinterAgent::start_print(PrintParams params, OnUpdateStatusFn update_fn, WasCancelledFn cancel_fn, OnWaitFn wait_fn)
|
||||
{
|
||||
// Moonraker has no cloud print path; report honestly instead of a false success
|
||||
// (the LAN path goes through start_local_print).
|
||||
(void) params;
|
||||
(void) update_fn;
|
||||
(void) cancel_fn;
|
||||
(void) wait_fn;
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
|
||||
int MoonrakerPrinterAgent::start_local_print_with_record(PrintParams params,
|
||||
@@ -579,11 +581,12 @@ int MoonrakerPrinterAgent::start_local_print_with_record(PrintParams params
|
||||
WasCancelledFn cancel_fn,
|
||||
OnWaitFn wait_fn)
|
||||
{
|
||||
// No record/cloud variant for Moonraker; the LAN path goes through start_local_print.
|
||||
(void) params;
|
||||
(void) update_fn;
|
||||
(void) cancel_fn;
|
||||
(void) wait_fn;
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
|
||||
int MoonrakerPrinterAgent::start_send_gcode_to_sdcard(PrintParams params,
|
||||
@@ -695,10 +698,12 @@ int MoonrakerPrinterAgent::start_local_print(PrintParams params, OnUpdateStatusF
|
||||
|
||||
int MoonrakerPrinterAgent::start_sdcard_print(PrintParams params, OnUpdateStatusFn update_fn, WasCancelledFn cancel_fn)
|
||||
{
|
||||
// Printing a file straight from the printer's storage is not implemented; say so
|
||||
// instead of reporting success for a no-op.
|
||||
(void) params;
|
||||
(void) update_fn;
|
||||
(void) cancel_fn;
|
||||
return BAMBU_NETWORK_SUCCESS;
|
||||
return ORCA_NETWORK_ERR_CMD_NOT_SUPPORTED;
|
||||
}
|
||||
|
||||
int MoonrakerPrinterAgent::set_on_ssdp_msg_fn(OnMsgArrivedFn fn)
|
||||
@@ -1951,49 +1956,59 @@ bool MoonrakerPrinterAgent::fetch_webcam_info(const ConnectionSettings& connecti
|
||||
std::string webcam_name;
|
||||
CameraStreamMode stream_mode = CameraStreamMode::none;
|
||||
std::string error;
|
||||
// A subclass may name its stream directly (e.g. a fixed webcam path with no
|
||||
// /server/webcams/list entry); consult it before doing any HTTP.
|
||||
const std::string override_url = webcam_stream_override(connection.base_url);
|
||||
try {
|
||||
std::string response_body;
|
||||
bool success = false;
|
||||
std::string http_error;
|
||||
|
||||
auto http = Http::get(join_url(connection.base_url, "/server/webcams/list"));
|
||||
configure_http(http, connection);
|
||||
if (!connection.api_key.empty()) {
|
||||
http.header("X-Api-Key", connection.api_key);
|
||||
}
|
||||
http.timeout_connect(5)
|
||||
.timeout_max(10)
|
||||
.on_complete([&](std::string body, unsigned status_code) {
|
||||
if (status_code == 200) {
|
||||
response_body = body;
|
||||
success = true;
|
||||
} else {
|
||||
http_error = "HTTP error: " + std::to_string(status_code);
|
||||
}
|
||||
})
|
||||
.on_error([&](std::string body, std::string err, unsigned status_code) {
|
||||
http_error = err;
|
||||
if (status_code > 0) {
|
||||
http_error += " (HTTP " + std::to_string(status_code) + ")";
|
||||
}
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
if (!success) {
|
||||
error = http_error.empty() ? "Connection failed" : http_error;
|
||||
if (!override_url.empty()) {
|
||||
camera_url = override_url;
|
||||
stream_mode = (override_url.rfind("rtsp://", 0) == 0 || override_url.rfind("rtsps://", 0) == 0)
|
||||
? CameraStreamMode::rtsp
|
||||
: CameraStreamMode::http;
|
||||
} else {
|
||||
BOOST_LOG_TRIVIAL(info) << "[Moonraker Diagnostic] " << connection.base_url << ":" << response_body;
|
||||
auto json = nlohmann::json::parse(response_body, nullptr, false, true);
|
||||
if (json.is_discarded()) {
|
||||
error = "Invalid JSON response";
|
||||
std::string response_body;
|
||||
bool success = false;
|
||||
std::string http_error;
|
||||
|
||||
auto http = Http::get(join_url(connection.base_url, "/server/webcams/list"));
|
||||
configure_http(http, connection);
|
||||
if (!connection.api_key.empty()) {
|
||||
http.header("X-Api-Key", connection.api_key);
|
||||
}
|
||||
http.timeout_connect(5)
|
||||
.timeout_max(10)
|
||||
.on_complete([&](std::string body, unsigned status_code) {
|
||||
if (status_code == 200) {
|
||||
response_body = body;
|
||||
success = true;
|
||||
} else {
|
||||
http_error = "HTTP error: " + std::to_string(status_code);
|
||||
}
|
||||
})
|
||||
.on_error([&](std::string body, std::string err, unsigned status_code) {
|
||||
http_error = err;
|
||||
if (status_code > 0) {
|
||||
http_error += " (HTTP " + std::to_string(status_code) + ")";
|
||||
}
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
if (!success) {
|
||||
error = http_error.empty() ? "Connection failed" : http_error;
|
||||
} else {
|
||||
MoonrakerWebcamSelection selection;
|
||||
if (moonraker_parse_webcam_list(json, connection.base_url, selection)) {
|
||||
camera_url = selection.url;
|
||||
stream_mode = selection.mode;
|
||||
webcam_name = selection.name;
|
||||
BOOST_LOG_TRIVIAL(info) << "[Moonraker Diagnostic] " << connection.base_url << ":" << response_body;
|
||||
auto json = nlohmann::json::parse(response_body, nullptr, false, true);
|
||||
if (json.is_discarded()) {
|
||||
error = "Invalid JSON response";
|
||||
} else {
|
||||
error = selection.error;
|
||||
MoonrakerWebcamSelection selection;
|
||||
if (moonraker_parse_webcam_list(json, connection.base_url, selection)) {
|
||||
camera_url = selection.url;
|
||||
stream_mode = selection.mode;
|
||||
webcam_name = selection.name;
|
||||
} else {
|
||||
error = selection.error;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2033,7 +2048,7 @@ bool MoonrakerPrinterAgent::post_print_action(const std::string& action,
|
||||
// /printer/gcode/script waits behind the gcode queue and can no-op while the
|
||||
// printer is busy (long move, heating, inside a macro).
|
||||
// note: empty JSON body avoids a body-less POST (curl would treat it as a
|
||||
// streamed upload) - same reason start_print_file sends a body.
|
||||
// streamed upload and fail).
|
||||
const std::string full_url = join_url(connection.base_url, "/printer/print/" + action);
|
||||
bool success = false;
|
||||
std::string http_error;
|
||||
@@ -3157,84 +3172,6 @@ int MoonrakerPrinterAgent::cancel_print(const std::string& dev_id)
|
||||
return post_print_action("cancel") ? BAMBU_NETWORK_SUCCESS : BAMBU_NETWORK_ERR_SEND_MSG_FAILED;
|
||||
}
|
||||
|
||||
bool MoonrakerPrinterAgent::start_print_file(const ConnectionSettings& connection,
|
||||
const std::string& filename,
|
||||
std::string& error_msg) const
|
||||
{
|
||||
// Start the given file (path relative to the gcodes root). The filename is
|
||||
// sent both as a query parameter and in the JSON body: Moonraker accepts
|
||||
// either, and sending a body avoids a body-less POST (which curl would treat
|
||||
// as a streamed upload and try to read via the file-read callback).
|
||||
std::string url = join_url(connection.base_url, "/printer/print/start") +
|
||||
"?filename=" + Http::url_encode(filename);
|
||||
|
||||
nlohmann::json payload;
|
||||
payload["filename"] = filename;
|
||||
|
||||
bool success = false;
|
||||
|
||||
auto http = Http::post(url);
|
||||
configure_http(http, connection);
|
||||
if (!connection.api_key.empty()) {
|
||||
http.header("X-Api-Key", connection.api_key);
|
||||
}
|
||||
http.header("Content-Type", "application/json")
|
||||
.set_post_body(payload.dump())
|
||||
.timeout_connect(5)
|
||||
.timeout_max(10)
|
||||
.on_complete([&](std::string body, unsigned status) {
|
||||
(void) body;
|
||||
if (status == 200) {
|
||||
success = true;
|
||||
} else {
|
||||
error_msg = "HTTP " + std::to_string(status);
|
||||
}
|
||||
})
|
||||
.on_error([&](std::string body, std::string err, unsigned status) {
|
||||
(void) body;
|
||||
error_msg = err;
|
||||
if (status > 0) {
|
||||
error_msg += " (HTTP " + std::to_string(status) + ")";
|
||||
}
|
||||
})
|
||||
.perform_sync();
|
||||
|
||||
if (success) {
|
||||
return true;
|
||||
}
|
||||
|
||||
// Moonraker holds the /printer/print/start response until the print actually
|
||||
// begins, so a slow PRINT_START (heating, homing, bed mesh) can exceed our HTTP
|
||||
// timeout even though the command was accepted and the print is starting. Don't
|
||||
// report failure on the HTTP result alone: poll print_stats and treat a
|
||||
// printing/paused state as success. print_stats.state flips to "printing" as soon
|
||||
// as the file starts streaming, which is earlier than the held HTTP response.
|
||||
BOOST_LOG_TRIVIAL(warning) << "MoonrakerPrinterAgent: start print not confirmed over HTTP (" << error_msg
|
||||
<< "); verifying print_stats state";
|
||||
for (int attempt = 0; attempt < 10; ++attempt) {
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(1500));
|
||||
|
||||
nlohmann::json status;
|
||||
std::string query_err;
|
||||
if (!query_printer_status(connection, status, query_err)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
std::string state;
|
||||
if (status.contains("print_stats") && status["print_stats"].contains("state") &&
|
||||
status["print_stats"]["state"].is_string()) {
|
||||
state = status["print_stats"]["state"].get<std::string>();
|
||||
}
|
||||
if (state == "printing" || state == "paused") {
|
||||
BOOST_LOG_TRIVIAL(info) << "MoonrakerPrinterAgent: print confirmed started (print_stats.state=" << state << ")";
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
BOOST_LOG_TRIVIAL(error) << "MoonrakerPrinterAgent: start print failed: " << error_msg << ", url: " << url;
|
||||
return false;
|
||||
}
|
||||
|
||||
void MoonrakerPrinterAgent::perform_connection_async(const std::string& dev_id,
|
||||
ConnectionSettings connection,
|
||||
uint64_t generation)
|
||||
|
||||
@@ -267,11 +267,6 @@ private:
|
||||
const ConnectionSettings& connection,
|
||||
OnUpdateStatusFn update_fn, WasCancelledFn cancel_fn);
|
||||
|
||||
// Start a print of a previously uploaded G-code file (path relative to the
|
||||
// Moonraker gcodes root).
|
||||
bool start_print_file(const ConnectionSettings& connection,
|
||||
const std::string& filename, std::string& error_msg) const;
|
||||
|
||||
// Connection thread management
|
||||
void perform_connection_async(const std::string& dev_id,
|
||||
ConnectionSettings connection,
|
||||
|
||||
@@ -101,6 +101,8 @@ bool QidiPrinterAgent::fetch_filament_info(std::string dev_id, FilamentSyncMode
|
||||
std::lock_guard<std::mutex> lock(fetch_lifecycle_mutex);
|
||||
if (shutting_down.load())
|
||||
return false;
|
||||
if (filament_fetch_in_flight.load() > 0)
|
||||
return true; // a fetch is already running; don't pile on
|
||||
filament_fetch_in_flight.fetch_add(1, std::memory_order_relaxed);
|
||||
}
|
||||
|
||||
@@ -157,43 +159,41 @@ bool QidiPrinterAgent::apply_box_mapping(const PrintParams& params) const
|
||||
// switch this gate to HasAms()/box_count instead.)
|
||||
const int enable = params.task_use_ams ? 1 : 0;
|
||||
const std::string dev_id = get_connection_settings().dev_id;
|
||||
if (!send_gcode(dev_id, "SAVE_VARIABLE VARIABLE=enable_box VALUE=" + std::to_string(enable))) {
|
||||
BOOST_LOG_TRIVIAL(error) << "QidiPrinterAgent::apply_box_mapping: failed to set enable_box";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Build one gcode/script request instead of N blocking HTTP calls: apply_box_mapping
|
||||
// runs on the caller's (GUI) thread before the print starts, so per-tool round trips
|
||||
// would freeze the UI.
|
||||
std::string script = "SAVE_VARIABLE VARIABLE=enable_box VALUE=" + std::to_string(enable);
|
||||
|
||||
// When the box isn't used this job, leave the existing value_t<tool> slot
|
||||
// assignments untouched (enable_box=0 is enough to disengage it).
|
||||
if (!enable)
|
||||
return true;
|
||||
|
||||
if (params.ams_mapping.empty()) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "QidiPrinterAgent::apply_box_mapping: enable_box set but ams_mapping is empty";
|
||||
return true;
|
||||
}
|
||||
|
||||
// ams_mapping (v0) is a JSON array indexed by filament/tool; each value is the
|
||||
// physical box slot (-1 = unmapped). Mirror it onto the printer's value_t<tool>
|
||||
// variables: SAVE_VARIABLE VARIABLE=value_t<tool> VALUE='slot<n>'.
|
||||
auto mapping = nlohmann::json::parse(params.ams_mapping, nullptr, /*allow_exceptions*/ false);
|
||||
if (mapping.is_discarded() || !mapping.is_array()) {
|
||||
BOOST_LOG_TRIVIAL(error) << "QidiPrinterAgent::apply_box_mapping: invalid ams_mapping: " << params.ams_mapping;
|
||||
return false;
|
||||
}
|
||||
|
||||
for (size_t tool = 0; tool < mapping.size(); ++tool) {
|
||||
if (!mapping[tool].is_number_integer())
|
||||
continue;
|
||||
const int slot = mapping[tool].get<int>();
|
||||
if (slot < 0)
|
||||
continue; // unmapped filament — skip
|
||||
const std::string gcode = "SAVE_VARIABLE VARIABLE=value_t" + std::to_string(tool) +
|
||||
" VALUE=\"'slot" + std::to_string(slot) + "'\"";
|
||||
if (!send_gcode(dev_id, gcode)) {
|
||||
BOOST_LOG_TRIVIAL(error) << "QidiPrinterAgent::apply_box_mapping: failed to set value_t" << tool;
|
||||
return false;
|
||||
if (enable) {
|
||||
if (params.ams_mapping.empty()) {
|
||||
BOOST_LOG_TRIVIAL(warning) << "QidiPrinterAgent::apply_box_mapping: enable_box set but ams_mapping is empty";
|
||||
} else {
|
||||
// ams_mapping (v0) is a JSON array indexed by filament/tool; each value is the
|
||||
// physical box slot (-1 = unmapped). Mirror it onto value_t<tool>.
|
||||
auto mapping = nlohmann::json::parse(params.ams_mapping, nullptr, /*allow_exceptions*/ false);
|
||||
if (mapping.is_discarded() || !mapping.is_array()) {
|
||||
BOOST_LOG_TRIVIAL(error) << "QidiPrinterAgent::apply_box_mapping: invalid ams_mapping: " << params.ams_mapping;
|
||||
return false;
|
||||
}
|
||||
for (size_t tool = 0; tool < mapping.size(); ++tool) {
|
||||
if (!mapping[tool].is_number_integer())
|
||||
continue;
|
||||
const int slot = mapping[tool].get<int>();
|
||||
if (slot < 0)
|
||||
continue; // unmapped filament — skip
|
||||
script += "\nSAVE_VARIABLE VARIABLE=value_t" + std::to_string(tool) +
|
||||
" VALUE=\"'slot" + std::to_string(slot) + "'\"";
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (!send_gcode(dev_id, script)) {
|
||||
BOOST_LOG_TRIVIAL(error) << "QidiPrinterAgent::apply_box_mapping: failed to send box mapping";
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -213,8 +213,9 @@ int QidiPrinterAgent::start_print(PrintParams params, OnUpdateStatusFn update_fn
|
||||
|
||||
int QidiPrinterAgent::start_local_print_with_record(PrintParams params, OnUpdateStatusFn update_fn, WasCancelledFn cancel_fn, OnWaitFn wait_fn)
|
||||
{
|
||||
// A failed box mapping is a send failure, not an upload failure.
|
||||
if (!apply_box_mapping(params))
|
||||
return BAMBU_NETWORK_ERR_PRINT_WR_UPLOAD_FTP_FAILED;
|
||||
return BAMBU_NETWORK_ERR_PRINT_LP_PUBLISH_MSG_FAILED;
|
||||
return MoonrakerPrinterAgent::start_local_print_with_record(std::move(params), update_fn, cancel_fn, wait_fn);
|
||||
}
|
||||
|
||||
|
||||
@@ -231,6 +231,8 @@ bool SnapmakerPrinterAgent::fetch_filament_info(std::string dev_id, FilamentSync
|
||||
std::lock_guard<std::mutex> lock(fetch_lifecycle_mutex);
|
||||
if (shutting_down.load())
|
||||
return false;
|
||||
if (filament_fetch_in_flight.load() > 0)
|
||||
return true; // a fetch is already running; don't pile on
|
||||
filament_fetch_in_flight.fetch_add(1, std::memory_order_relaxed);
|
||||
}
|
||||
|
||||
|
||||
@@ -31,6 +31,33 @@ namespace Slic3r { class IPrinterAgent; }
|
||||
using namespace Slic3r;
|
||||
namespace py = pybind11;
|
||||
|
||||
namespace {
|
||||
|
||||
// Releases a promise on scope exit, so a throwing REQUIRE cannot leave a parked detached
|
||||
// thread (and any destructor that joins it) blocked forever.
|
||||
class ScopedPromiseRelease
|
||||
{
|
||||
public:
|
||||
explicit ScopedPromiseRelease(std::shared_ptr<std::promise<void>> p) : m_p(std::move(p)) {}
|
||||
~ScopedPromiseRelease()
|
||||
{
|
||||
if (m_p) {
|
||||
try {
|
||||
m_p->set_value();
|
||||
} catch (...) {
|
||||
// promise already satisfied
|
||||
}
|
||||
}
|
||||
}
|
||||
ScopedPromiseRelease(const ScopedPromiseRelease&) = delete;
|
||||
ScopedPromiseRelease& operator=(const ScopedPromiseRelease&) = delete;
|
||||
|
||||
private:
|
||||
std::shared_ptr<std::promise<void>> m_p;
|
||||
};
|
||||
|
||||
} // namespace
|
||||
|
||||
class MoonrakerParserProbe : public MoonrakerPrinterAgent
|
||||
{
|
||||
public:
|
||||
@@ -202,25 +229,31 @@ TEST_CASE("unit: a fire-and-forget override of fetch_filament_info is not waited
|
||||
public:
|
||||
explicit RecordingAgent(std::string log_dir) : MoonrakerPrinterAgent(std::move(log_dir)) {}
|
||||
|
||||
std::atomic<bool> invoked{false};
|
||||
std::promise<void> release_gate;
|
||||
std::promise<void> done_promise;
|
||||
// Shared so the detached proxy fetch never touches `this`: a throwing REQUIRE
|
||||
// then cannot leave it dereferencing a destroyed agent.
|
||||
std::shared_ptr<std::atomic<bool>> invoked{std::make_shared<std::atomic<bool>>(false)};
|
||||
std::shared_ptr<std::promise<void>> release_gate{std::make_shared<std::promise<void>>()};
|
||||
std::shared_ptr<std::promise<void>> done_promise{std::make_shared<std::promise<void>>()};
|
||||
|
||||
bool fetch_filament_info(std::string /*dev_id*/, FilamentSyncMode /*sync_mode*/ = FilamentSyncMode::pull) override
|
||||
{
|
||||
std::thread([this]() {
|
||||
invoked.store(true);
|
||||
auto invoked_p = invoked;
|
||||
auto release_gate_p = release_gate;
|
||||
auto done_promise_p = done_promise;
|
||||
std::thread([invoked_p, release_gate_p, done_promise_p]() {
|
||||
invoked_p->store(true);
|
||||
// Block here until the test explicitly releases us, proving the caller
|
||||
// (fetch_filament_info) does not wait for this to run.
|
||||
release_gate.get_future().wait();
|
||||
done_promise.set_value();
|
||||
release_gate_p->get_future().wait();
|
||||
done_promise_p->set_value();
|
||||
}).detach();
|
||||
return true;
|
||||
}
|
||||
};
|
||||
|
||||
auto agent = std::make_shared<RecordingAgent>(std::string{});
|
||||
auto done_future = agent->done_promise.get_future();
|
||||
auto done_future = agent->done_promise->get_future();
|
||||
ScopedPromiseRelease release_gate_guard{agent->release_gate};
|
||||
|
||||
bool immediate_result = agent->fetch_filament_info("test-dev");
|
||||
|
||||
@@ -230,9 +263,9 @@ TEST_CASE("unit: a fire-and-forget override of fetch_filament_info is not waited
|
||||
REQUIRE(done_future.wait_for(std::chrono::milliseconds(100)) == std::future_status::timeout);
|
||||
|
||||
// Now let the background call finish and confirm it actually ran (polymorphic dispatch).
|
||||
agent->release_gate.set_value();
|
||||
agent->release_gate->set_value();
|
||||
REQUIRE(done_future.wait_for(std::chrono::seconds(2)) == std::future_status::ready);
|
||||
REQUIRE(agent->invoked.load() == true);
|
||||
REQUIRE(agent->invoked->load() == true);
|
||||
}
|
||||
|
||||
namespace {
|
||||
@@ -241,17 +274,33 @@ namespace {
|
||||
std::atomic<int> g_deferred_fetch_running{0};
|
||||
std::atomic<bool> g_deferred_destroy_returned{false};
|
||||
|
||||
// Joins on scope exit so a throwing REQUIRE does not std::terminate.
|
||||
// Releases the given gates, then joins on scope exit: so a throwing REQUIRE cannot leave
|
||||
// the thread blocked (deadlocking the join) or let it std::terminate.
|
||||
class ScopedJoiner
|
||||
{
|
||||
public:
|
||||
explicit ScopedJoiner(std::thread& t) : m_thread(t) {}
|
||||
~ScopedJoiner() { if (m_thread.joinable()) m_thread.join(); }
|
||||
ScopedJoiner(std::thread& t, std::shared_ptr<std::promise<void>> gate1, std::shared_ptr<std::promise<void>> gate2)
|
||||
: m_thread(t), m_gates{std::move(gate1), std::move(gate2)}
|
||||
{}
|
||||
~ScopedJoiner()
|
||||
{
|
||||
for (auto& gate : m_gates) {
|
||||
if (gate) {
|
||||
try {
|
||||
gate->set_value();
|
||||
} catch (...) {
|
||||
// promise already satisfied
|
||||
}
|
||||
}
|
||||
}
|
||||
if (m_thread.joinable()) m_thread.join();
|
||||
}
|
||||
ScopedJoiner(const ScopedJoiner&) = delete;
|
||||
ScopedJoiner& operator=(const ScopedJoiner&) = delete;
|
||||
|
||||
private:
|
||||
std::thread& m_thread;
|
||||
std::thread& m_thread;
|
||||
std::shared_ptr<std::promise<void>> m_gates[2];
|
||||
};
|
||||
|
||||
// A fetch that parks before touching the in-flight counter, so teardown's wait can
|
||||
@@ -265,6 +314,7 @@ public:
|
||||
std::shared_ptr<std::promise<void>> entered{std::make_shared<std::promise<void>>()};
|
||||
std::shared_ptr<std::promise<void>> allow_fetch{std::make_shared<std::promise<void>>()};
|
||||
std::shared_ptr<std::promise<void>> allow_finish{std::make_shared<std::promise<void>>()};
|
||||
std::shared_ptr<std::promise<void>> running{std::make_shared<std::promise<void>>()};
|
||||
|
||||
// Runs the callable on the command worker, which teardown joins.
|
||||
void post(std::function<void()> fn) { enqueue_command(std::move(fn)); }
|
||||
@@ -275,12 +325,13 @@ public:
|
||||
auto entered_p = entered;
|
||||
auto allow_fetch_p = allow_fetch;
|
||||
auto allow_finish_p = allow_finish;
|
||||
auto running_p = running;
|
||||
|
||||
entered_p->set_value();
|
||||
allow_fetch_p->get_future().wait();
|
||||
|
||||
filament_fetch_in_flight.fetch_add(1, std::memory_order_relaxed);
|
||||
std::thread([this, finish = std::move(allow_finish_p)] {
|
||||
std::thread([this, finish = std::move(allow_finish_p), running = std::move(running_p)] {
|
||||
struct InFlightGuard
|
||||
{
|
||||
MoonrakerPrinterAgent& owner;
|
||||
@@ -288,6 +339,7 @@ public:
|
||||
} guard{*this};
|
||||
|
||||
g_deferred_fetch_running.fetch_add(1, std::memory_order_relaxed);
|
||||
running->set_value();
|
||||
finish->get_future().wait();
|
||||
g_deferred_fetch_running.fetch_sub(1, std::memory_order_relaxed);
|
||||
}).detach();
|
||||
@@ -310,6 +362,12 @@ TEST_CASE("an agent's destruction waits for a fetch started by its worker during
|
||||
auto entered = agent->entered;
|
||||
auto allow_fetch = agent->allow_fetch;
|
||||
auto allow_finish = agent->allow_finish;
|
||||
auto running = agent->running;
|
||||
|
||||
// Safety net for the pre-destroyer failure paths: release both gates before the
|
||||
// agent is destroyed (declared after it, so destroyed before it).
|
||||
ScopedPromiseRelease release_finish{allow_finish};
|
||||
ScopedPromiseRelease release_fetch{allow_fetch};
|
||||
|
||||
// Park a fetch inside the command worker while the agent is still complete.
|
||||
agent->post([ptr = agent.get()] { ptr->fetch_filament_info("dev", FilamentSyncMode::pull); });
|
||||
@@ -320,15 +378,13 @@ TEST_CASE("an agent's destruction waits for a fetch started by its worker during
|
||||
owned.reset();
|
||||
g_deferred_destroy_returned.store(true);
|
||||
});
|
||||
ScopedJoiner join_destroyer{destroyer};
|
||||
// Releases both gates before joining, so a failing REQUIRE cannot deadlock the join.
|
||||
ScopedJoiner join_destroyer{destroyer, allow_fetch, allow_finish};
|
||||
|
||||
// Let teardown pass its wait; the worker has not reserved yet.
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(300));
|
||||
// Let the worker reserve the in-flight slot and spawn its fetch, then wait until it
|
||||
// is genuinely parked (no polling).
|
||||
allow_fetch->set_value();
|
||||
|
||||
// Get the fetch actually in flight (parked on allow_finish).
|
||||
for (int i = 0; i < 200 && g_deferred_fetch_running.load() == 0; ++i)
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(10));
|
||||
REQUIRE(running->get_future().wait_for(std::chrono::seconds(5)) == std::future_status::ready);
|
||||
REQUIRE(g_deferred_fetch_running.load() == 1);
|
||||
|
||||
// A correct teardown cannot return while the fetch is parked; give a buggy one time.
|
||||
|
||||
Reference in New Issue
Block a user