diff --git a/src/slic3r/GUI/GUI_App.cpp b/src/slic3r/GUI/GUI_App.cpp index 756775a4d9..9f7202d8c0 100644 --- a/src/slic3r/GUI/GUI_App.cpp +++ b/src/slic3r/GUI/GUI_App.cpp @@ -1861,6 +1861,17 @@ namespace { constexpr int FINAL_DRAIN_TIMEOUT_MS = 100; // Final event processing before destruction constexpr int POLL_INTERVAL_MS = 50; // Polling interval for state checks constexpr int MAX_YIELD_ITERATIONS = 20; // Maximum wxYield calls per drain cycle + + // The Orca agent's health check reports into GUI_App from a worker thread, and the agent can + // outlive GUI_App's use of it (the plugin service keeps a reference), so stop the check before + // the agent is deleted or replaced. The cast is null when no Orca agent is attached. + void stop_orca_health_check(NetworkAgent* agent) + { + if (!agent) + return; + if (auto orca_agent = std::dynamic_pointer_cast(agent->get_cloud_agent())) + orca_agent->stop_health_check(); + } } // Process pending wx events with bounded iteration count @@ -1993,6 +2004,10 @@ bool GUI_App::hot_reload_network_plugin() } if (m_agent) { + // Stop the health check first: a result it posts after the Phase 2 drain would run while + // m_agent is null or replaced. + stop_orca_health_check(m_agent); + // Phase 1: Clear all callbacks (stops new invocations) BOOST_LOG_TRIVIAL(info) << __FUNCTION__ << ": Phase 1 - clearing callbacks"; m_agent->set_on_ssdp_msg_fn(nullptr); @@ -2860,6 +2875,7 @@ int GUI_App::OnExit() NetworkAgentFactory::clear_printer_agent_cache(); if (m_agent) { + stop_orca_health_check(m_agent); // BBS avoid a crash on mac platform #ifdef __WINDOWS__ m_agent->start_discovery(false, false); @@ -2893,6 +2909,14 @@ int GUI_App::OnExit() return wxApp::OnExit(); } +void GUI_App::CleanUp() +{ + // OnExit() does not run when OnInit() fails, and m_agent is then still set. Stop the check before + // wxApp::CleanUp() resets wxTheApp, which a result posted from the worker would still be using. + stop_orca_health_check(m_agent); + wxApp::CleanUp(); +} + class wxBoostLog : public wxLog { void DoLogText(const wxString &msg) override { @@ -3115,9 +3139,11 @@ bool GUI_App::on_init_inner() // OnExit() and ~GUI_App() never run. Shut the plugins and Python down here as ~GUI_App() does. Left to // PluginManager's static destructor, the shutdown locks hook state that has already been destroyed and aborts. // Unload the Bambu network plugin too. Its static destructors abort if its agent's threads are still running. + // Stop the Orca Cloud health check as OnExit() does, so no request is in flight while exit() tears down. wxGetApp().Bind(wxEVT_END_SESSION, [this](wxCloseEvent &e) { BOOST_LOG_TRIVIAL(info) << __FUNCTION__ << "received wxEVT_END_SESSION"; stop_sync_user_preset(); + stop_orca_health_check(m_agent); Slic3r::NetworkAgent::unload_network_module(); Slic3r::PluginManager::instance().shutdown(); Slic3r::PythonInterpreter::instance().shutdown(); @@ -3975,6 +4001,9 @@ bool GUI_App::on_init_network(bool try_backup) // m_agent = new Slic3r::NetworkAgent(data_directory); std::unique_ptr agent_ptr = Slic3r::create_agent_from_config(data_directory, app_config); + // A direct restart_networking() replaces m_agent without deleting it, so a health check in flight + // on the old agent would still report into the GUI. + stop_orca_health_check(m_agent); m_agent = agent_ptr.release(); if (!m_device_manager) diff --git a/src/slic3r/GUI/GUI_App.hpp b/src/slic3r/GUI/GUI_App.hpp index 361693ed42..8da0f5ea33 100644 --- a/src/slic3r/GUI/GUI_App.hpp +++ b/src/slic3r/GUI/GUI_App.hpp @@ -373,6 +373,7 @@ public: std::string get_local_models_path(); bool OnInit() override; int OnExit() override; + void CleanUp() override; bool initialized() const { return m_initialized; } inline bool is_enable_multi_machine() { return this->app_config&& this->app_config->get("enable_multi_machine") == "true"; } #ifdef SLIC3R_CAD diff --git a/src/slic3r/Utils/ICloudServiceAgent.hpp b/src/slic3r/Utils/ICloudServiceAgent.hpp index 6457b3fda6..1bcb4d370a 100644 --- a/src/slic3r/Utils/ICloudServiceAgent.hpp +++ b/src/slic3r/Utils/ICloudServiceAgent.hpp @@ -200,6 +200,9 @@ public: /** * Force a server state recheck, clearing any cached state. + * The result arrives through is_server_connected() and OnServerConnectedFn. Implementations may + * return before the check finishes, fold the call into a check already in flight, and ignore + * calls while the agent is shutting down. */ virtual int refresh_connection() = 0; diff --git a/src/slic3r/Utils/OrcaCloudServiceAgent.cpp b/src/slic3r/Utils/OrcaCloudServiceAgent.cpp index 7aa2b1ab2e..5a5cedeb87 100644 --- a/src/slic3r/Utils/OrcaCloudServiceAgent.cpp +++ b/src/slic3r/Utils/OrcaCloudServiceAgent.cpp @@ -523,6 +523,7 @@ OrcaCloudServiceAgent::OrcaCloudServiceAgent(std::string log_dir) OrcaCloudServiceAgent::~OrcaCloudServiceAgent() { + stop_health_check(); if (refresh_thread.joinable()) { refresh_thread.join(); } @@ -963,20 +964,55 @@ bool OrcaCloudServiceAgent::ensure_token_fresh(const std::string& reason) { retu // ICloudServiceAgent - Server Connectivity // ============================================================================ -int OrcaCloudServiceAgent::connect_server() +// /api/v1/health needs no auth, so unlike the data requests this does no token refresh or 401 +// recovery: refresh_connection() runs it off the UI thread, where a refresh could race a logout, and +// it has to stay cancellable. +bool OrcaCloudServiceAgent::run_health_check(const std::atomic_bool* cancel) { - std::string response; - unsigned int http_code = 0; - int result = http_get(ORCA_HEALTH_PATH, &response, &http_code); + if (cancel && cancel->load()) + return false; - bool connected = (result == BAMBU_NETWORK_SUCCESS && http_code >= 200 && http_code < 300); + HttpResult res; + try { + auto http = Http::get(api_base_url + ORCA_HEALTH_PATH); + http.tls_verify(true); + for (const auto& [key, value] : data_headers()) + http.header(key, value); + if (cancel) + http.on_progress([cancel](Http::Progress, bool& cancel_request) { cancel_request = cancel->load(); }); + http.on_complete([&](std::string body, unsigned status) { + res.success = true; + res.status = status; + res.body = std::move(body); + }) + .on_error([&](std::string body, std::string error, unsigned status) { + res.status = status == 0 ? 404 : status; // same mapping as http_get + res.body = std::move(body); + BOOST_LOG_TRIVIAL(error) << "OrcaCloudServiceAgent: health check failed - " << error; + }) + .timeout_max(30) + .perform_sync(); + } catch (const std::exception& e) { + BOOST_LOG_TRIVIAL(error) << "OrcaCloudServiceAgent: health check exception - " << e.what(); + } + // Stopped mid-request: report nothing, the receivers may be going away. + if (cancel && cancel->load()) + return false; + + const bool connected = res.success; { std::lock_guard lock(state_mutex); is_connected = connected; } + if (!connected) + invoke_http_error_callback(res.status, res.body); + invoke_server_connected_callback(connected ? 0 : -1, res.status); + return connected; +} - invoke_server_connected_callback(connected ? 0 : -1, http_code); - return connected ? BAMBU_NETWORK_SUCCESS : BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED; +int OrcaCloudServiceAgent::connect_server() +{ + return run_health_check(nullptr) ? BAMBU_NETWORK_SUCCESS : BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED; } bool OrcaCloudServiceAgent::is_server_connected() @@ -985,7 +1021,44 @@ bool OrcaCloudServiceAgent::is_server_connected() return is_connected; } -int OrcaCloudServiceAgent::refresh_connection() { return connect_server(); } +// The device manager calls this every 5 s on the UI thread and the GET can take up to its 30 s +// timeout, so the check runs on a worker and reports through is_server_connected() and the +// callbacks. A call while a check is in flight is folded into it. Called from one thread only (the +// UI thread), like stop_health_check(). +int OrcaCloudServiceAgent::refresh_connection() +{ + if (health_check_stopped.load()) + return BAMBU_NETWORK_ERR_CANCELED; + bool expected = false; + if (!health_check_running.compare_exchange_strong(expected, true)) + return BAMBU_NETWORK_SUCCESS; + if (health_check_thread.joinable()) + health_check_thread.join(); // the previous check is done; only its thread exit is left + try { + health_check_thread = std::thread([this] { + try { + run_health_check(&health_check_stopped); + } catch (...) { + // Callback dispatch (std::function copies, CallAfter) must not terminate the app. + BOOST_LOG_TRIVIAL(error) << "OrcaCloudServiceAgent: health check worker exception"; + } + health_check_running.store(false); + }); + } catch (...) { + // std::system_error, or std::bad_alloc for the callable: leave the next tick free to retry. + health_check_running.store(false); + BOOST_LOG_TRIVIAL(error) << "OrcaCloudServiceAgent: cannot start health check"; + return BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED; + } + return BAMBU_NETWORK_SUCCESS; +} + +void OrcaCloudServiceAgent::stop_health_check() +{ + health_check_stopped.store(true); + if (health_check_thread.joinable()) + health_check_thread.join(); +} int OrcaCloudServiceAgent::start_subscribe(std::string module) { diff --git a/src/slic3r/Utils/OrcaCloudServiceAgent.hpp b/src/slic3r/Utils/OrcaCloudServiceAgent.hpp index d4a182db12..ae9e1f3213 100644 --- a/src/slic3r/Utils/OrcaCloudServiceAgent.hpp +++ b/src/slic3r/Utils/OrcaCloudServiceAgent.hpp @@ -206,6 +206,10 @@ public: int connect_server() override; bool is_server_connected() override; int refresh_connection() override; + // Cancels a health check started by refresh_connection() and waits for it; later + // refresh_connection() calls do nothing. Call it before tearing down what the server-connected + // and HTTP-error callbacks reach. + void stop_health_check(); bool is_refresh_running() const { return refresh_running.load(); } int start_subscribe(std::string module) override; int stop_subscribe(std::string module) override; @@ -392,6 +396,10 @@ private: bool decode_jwt_expiry(const std::string& token, std::chrono::system_clock::time_point& out_tp); bool should_refresh_locked(std::chrono::seconds skew) const; + // Server reachability probe shared by connect_server() and refresh_connection(); returns + // false without reporting anything once `cancel` is set. + bool run_health_check(const std::atomic_bool* cancel); + // Callback invocation void invoke_server_connected_callback(int return_code, int reason_code); void invoke_http_error_callback(unsigned http_code, const std::string& http_body); @@ -454,6 +462,9 @@ private: mutable std::recursive_mutex state_mutex; std::thread refresh_thread; std::atomic_bool refresh_running{false}; + std::thread health_check_thread; + std::atomic_bool health_check_running{false}; + std::atomic_bool health_check_stopped{false}; }; } // namespace Slic3r diff --git a/tests/slic3rutils/test_orca_cloud_agent.cpp b/tests/slic3rutils/test_orca_cloud_agent.cpp index 8657310f27..99c4180a9a 100644 --- a/tests/slic3rutils/test_orca_cloud_agent.cpp +++ b/tests/slic3rutils/test_orca_cloud_agent.cpp @@ -2,14 +2,38 @@ #include #include +#include +#include +#include +#include +#include +#include +#include +#include #include #include +#include +#include +#include +#include +#include +#include +#include +#include #include +#include +#include #include +#include +#include +#include +#include #include +#include "slic3r/Utils/ICloudServiceAgent.hpp" #include "slic3r/Utils/OrcaCloudServiceAgent.hpp" +#include "slic3r/Utils/bambu_networking.hpp" #include "test_utils.hpp" using namespace Slic3r; @@ -61,6 +85,227 @@ std::string resolved_display_name(const nlohmann::json& session) return agent->get_user_nickname(); } +using namespace std::chrono_literals; +using boost::asio::ip::tcp; + +// A loopback HTTP server for the health check to talk to. It answers every request according to +// the current mode and records what it was sent. +class LoopbackServer +{ +public: + enum class Mode { Ok, Unavailable, Silent }; + + struct Request + { + std::string request_line; + bool has_authorization{false}; + }; + + LoopbackServer() : m_acceptor(m_io, tcp::endpoint(boost::asio::ip::make_address("127.0.0.1"), 0)) + { + accept_next(); + m_thread = std::thread([this] { m_io.run(); }); + } + + ~LoopbackServer() + { + m_io.stop(); + m_thread.join(); + } + + std::string url() const { return "http://127.0.0.1:" + std::to_string(m_acceptor.local_endpoint().port()); } + void set_mode(Mode mode) { m_mode = mode; } + + bool wait_for_accepts(size_t count, std::chrono::milliseconds timeout) + { + std::unique_lock lock(m_mutex); + return m_cv.wait_for(lock, timeout, [&] { return m_connections.size() >= count; }); + } + + std::vector requests() const + { + std::lock_guard lock(m_mutex); + return m_requests; + } + +private: + struct Connection + { + explicit Connection(tcp::socket socket) : socket(std::move(socket)) {} + tcp::socket socket; + boost::asio::streambuf buffer; + }; + + void accept_next() + { + m_acceptor.async_accept([this](boost::system::error_code ec, tcp::socket socket) { + if (ec) + return; + auto connection = std::make_shared(std::move(socket)); + { + std::lock_guard lock(m_mutex); + // Held here so a Silent connection stays open until the client gives up. + m_connections.push_back(connection); + } + m_cv.notify_all(); + read_request(connection); + accept_next(); + }); + } + + void read_request(const std::shared_ptr& connection) + { + boost::asio::async_read_until(connection->socket, connection->buffer, "\r\n\r\n", + [this, connection](boost::system::error_code ec, size_t) { + if (ec) + return; + std::istream in(&connection->buffer); + Request request; + std::string line; + std::getline(in, line); + request.request_line = line.substr(0, line.find('\r')); + while (std::getline(in, line) && line != "\r") + request.has_authorization |= boost::istarts_with(line, "authorization:"); + { + std::lock_guard lock(m_mutex); + m_requests.push_back(request); + } + + const Mode mode = m_mode; + if (mode == Mode::Silent) + return; + const std::string_view response = mode == Mode::Ok ? + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 11\r\nConnection: close\r\n\r\n{\"ok\":true}" : + "HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"; + boost::asio::async_write(connection->socket, boost::asio::buffer(response), + [connection](boost::system::error_code, size_t) { + boost::system::error_code ignored; + connection->socket.shutdown(tcp::socket::shutdown_both, ignored); + }); + }); + } + + // Declared first so it outlives the sockets that use it. + boost::asio::io_context m_io; + tcp::acceptor m_acceptor; + std::thread m_thread; + std::atomic m_mode{Mode::Ok}; + mutable std::mutex m_mutex; + std::condition_variable m_cv; + std::vector m_requests; + std::vector> m_connections; +}; + +// Sets an environment variable for the lifetime of the guard and restores the previous value. +class ScopedEnvVar +{ +public: + ScopedEnvVar(std::string name, const std::string& value) : m_name(std::move(name)) + { + if (const char* old = boost::nowide::getenv(m_name.c_str())) + m_previous = old; + boost::nowide::setenv(m_name.c_str(), value.c_str(), 1); + } + ~ScopedEnvVar() + { + if (m_previous) + boost::nowide::setenv(m_name.c_str(), m_previous->c_str(), 1); + else + boost::nowide::unsetenv(m_name.c_str()); + } + +private: + std::string m_name; + std::optional m_previous; +}; + +// What the agent reported through its callbacks, recorded on whichever thread delivers it so the +// assertions can run on the test thread. +struct CallbackLog +{ + void add_connected(int return_code, int reason_code) + { + { + std::lock_guard lock(mutex); + connected.emplace_back(return_code, reason_code); + } + cv.notify_all(); + } + + void add_http_error(unsigned http_code) + { + std::lock_guard lock(mutex); + http_errors.push_back(http_code); + } + + void add_queued(std::function task) + { + { + std::lock_guard lock(mutex); + queued.push_back(std::move(task)); + } + cv.notify_all(); + } + + bool wait_for_connected(size_t count, std::chrono::milliseconds timeout = 10s) + { + std::unique_lock lock(mutex); + return cv.wait_for(lock, timeout, [&] { return connected.size() >= count; }); + } + + bool wait_for_queued(size_t count, std::chrono::milliseconds timeout = 10s) + { + std::unique_lock lock(mutex); + return cv.wait_for(lock, timeout, [&] { return queued.size() >= count; }); + } + + std::mutex mutex; + std::condition_variable cv; + std::vector> connected; + std::vector http_errors; + std::vector> queued; +}; + +// An agent whose API and auth endpoints both point at a loopback server, with its callbacks +// recorded. Members are declared so the agent, and the health check thread it owns, is destroyed +// before the server, the log and the storage it uses. +struct HealthCheckFixture +{ + HealthCheckFixture() : agent(make_file_backed_agent(dir.path())) + { + agent->set_api_base_url(server.url()); + agent->set_auth_base_url(server.url()); + agent->set_on_server_connected_fn([this](CloudEvent, int return_code, int reason_code) { + log.add_connected(return_code, reason_code); + }); + agent->set_on_http_error_fn([this](CloudEvent, unsigned http_code, std::string) { log.add_http_error(http_code); }); + } + + // Defers callbacks the way the GUI's CallAfter does, instead of running them on the worker. + void queue_callbacks() { agent->set_queue_on_main_fn([this](std::function task) { log.add_queued(std::move(task)); }); } + + // Starts the next health check, retrying while a call is still folded into the previous one: + // its callback fires before the worker thread finishes. + bool start_check(size_t expected_accepts) + { + const auto deadline = std::chrono::steady_clock::now() + 5s; + do { + agent->refresh_connection(); + if (server.wait_for_accepts(expected_accepts, 10ms)) + return true; + } while (std::chrono::steady_clock::now() < deadline); + return false; + } + + // Loopback traffic must not go through a proxy configured in the environment. curl reads no_proxy + // before NO_PROXY. + ScopedEnvVar no_proxy{"no_proxy", "127.0.0.1"}; + LoopbackServer server; + CallbackLog log; + ScopedTemporaryDir dir{"orca-health"}; + std::unique_ptr agent; +}; + } // namespace TEST_CASE("Logging out removes the secret this instance saved", "[OrcaCloudServiceAgent]") @@ -174,3 +419,124 @@ TEST_CASE("Orca cloud nested session resolves display name consistently", "[Orca {"username", "orca_username"} })) == "orca_username"); } + +TEST_CASE_METHOD(HealthCheckFixture, "Refreshing the server connection returns before the health check finishes", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Silent); + + const auto start = std::chrono::steady_clock::now(); + CHECK(agent->refresh_connection() == BAMBU_NETWORK_SUCCESS); + // A synchronous check would wait for the request's 30 s timeout. + CHECK(std::chrono::steady_clock::now() - start < 2s); + + std::lock_guard lock(log.mutex); + CHECK(log.connected.empty()); +} + +TEST_CASE_METHOD(HealthCheckFixture, "Successive health checks follow the server's state", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Ok); + REQUIRE(start_check(1)); + REQUIRE(log.wait_for_connected(1)); + CHECK(agent->is_server_connected()); + + server.set_mode(LoopbackServer::Mode::Unavailable); + REQUIRE(start_check(2)); + REQUIRE(log.wait_for_connected(2)); + CHECK_FALSE(agent->is_server_connected()); + + server.set_mode(LoopbackServer::Mode::Ok); + REQUIRE(start_check(3)); + REQUIRE(log.wait_for_connected(3)); + CHECK(agent->is_server_connected()); + + std::lock_guard lock(log.mutex); + CHECK(log.connected == std::vector>{{0, 200}, {-1, 503}, {0, 200}}); + CHECK(log.http_errors == std::vector{503}); +} + +TEST_CASE_METHOD(HealthCheckFixture, "A refresh while a health check is in flight starts no second check", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Silent); + + agent->refresh_connection(); + REQUIRE(server.wait_for_accepts(1, 5s)); + + agent->refresh_connection(); + agent->refresh_connection(); + CHECK_FALSE(server.wait_for_accepts(2, 200ms)); +} + +TEST_CASE_METHOD(HealthCheckFixture, "Stopping the health check cancels the request in flight", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Silent); + queue_callbacks(); + + agent->refresh_connection(); + REQUIRE(server.wait_for_accepts(1, 5s)); + + const auto start = std::chrono::steady_clock::now(); + agent->stop_health_check(); + // Well under the request's 30 s timeout. + CHECK(std::chrono::steady_clock::now() - start < 5s); + CHECK(agent->refresh_connection() == BAMBU_NETWORK_ERR_CANCELED); + + std::lock_guard lock(log.mutex); + CHECK(log.connected.empty()); + CHECK(log.http_errors.empty()); + CHECK(log.queued.empty()); +} + +TEST_CASE_METHOD(HealthCheckFixture, "Connecting to the server reports the result synchronously", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Ok); + CHECK(agent->connect_server() == BAMBU_NETWORK_SUCCESS); + CHECK(agent->is_server_connected()); + + server.set_mode(LoopbackServer::Mode::Unavailable); + CHECK(agent->connect_server() == BAMBU_NETWORK_ERR_CONNECTION_TO_SERVER_FAILED); + CHECK_FALSE(agent->is_server_connected()); +} + +TEST_CASE_METHOD(HealthCheckFixture, "A health check result queued before stopping is still delivered and nothing is queued after", "[OrcaCloudServiceAgent]") +{ + server.set_mode(LoopbackServer::Mode::Ok); + queue_callbacks(); + + agent->refresh_connection(); + REQUIRE(log.wait_for_queued(1)); + agent->stop_health_check(); + + std::function task; + { + std::lock_guard lock(log.mutex); + REQUIRE(log.queued.size() == 1); + task = log.queued.front(); + } + // The GUI may run a task it queued after the agent is gone; the task must not need the agent. + agent.reset(); + task(); + + std::lock_guard lock(log.mutex); + CHECK(log.connected == std::vector>{{0, 200}}); +} + +TEST_CASE_METHOD(HealthCheckFixture, "The health check does no token work", "[OrcaCloudServiceAgent]") +{ + // A non-JWT access token has no expiry, so the session reads as due for a refresh. + REQUIRE(agent->set_user_session(flat_session_json({{"refresh_token", "test-refresh-token"}}), false)); + server.set_mode(LoopbackServer::Mode::Ok); + + REQUIRE(start_check(1)); + REQUIRE(log.wait_for_connected(1)); + CHECK(agent->connect_server() == BAMBU_NETWORK_SUCCESS); + + const auto requests = server.requests(); + REQUIRE(requests.size() == 2); + for (const auto& request : requests) { + CHECK(boost::starts_with(request.request_line, "GET /api/v1/health ")); + CHECK_FALSE(request.has_authorization); + } + CHECK(agent->is_user_login()); + CHECK(agent->get_access_token() == "test-token"); +}