diff --git a/dUgcServer/UgcProcessor.cpp b/dUgcServer/UgcProcessor.cpp index 00bee1f3b..d53b7d4f8 100644 --- a/dUgcServer/UgcProcessor.cpp +++ b/dUgcServer/UgcProcessor.cpp @@ -133,6 +133,7 @@ void UgcProcessor::Start() { } m_NextEviction = std::chrono::steady_clock::now(); m_Stopping = false; + UgcThrottle::Cancel(false); for (size_t i = 0; i < std::max(m_Config.threads, 1); i++) m_Threads.emplace_back(&UgcProcessor::Worker, this); LOG("UGC processing started with %zu worker(s), %llu MB stored", m_Threads.size(), static_cast(m_StoredBytes / (1024 * 1024))); } @@ -143,6 +144,9 @@ void UgcProcessor::Stop() { m_Stopping = true; m_Jobs.clear(); } + // Jobs being made stop at their next checkpoint instead of finishing (a big model can take minutes); their rows stay + // waiting and are made again after the restart + UgcThrottle::Cancel(true); m_Wake.notify_all(); for (auto& thread : m_Threads) { if (thread.joinable()) thread.join(); @@ -256,6 +260,10 @@ void UgcProcessor::Worker() { reply.contentType = eContentType::TEXT_PLAIN; reply.message = outcome.error; } + } catch (const UgcThrottle::Cancelled&) { + reply.status = eHTTPStatusCode::SERVICE_UNAVAILABLE; + reply.contentType = eContentType::TEXT_PLAIN; + reply.message = "the UGC server is stopping"; } catch (const std::exception& ex) { reply.status = eHTTPStatusCode::INTERNAL_SERVER_ERROR; reply.contentType = eContentType::TEXT_PLAIN; @@ -288,6 +296,12 @@ void UgcProcessor::Worker() { ? UgcJobs::ProcessModel(job.blob, m_Library, settings, static_cast(job.id), job.iconValues) : UgcJobs::ProcessModular(job.modular, m_Library.GetResPath(), settings); } + } catch (const UgcThrottle::Cancelled&) { + // Stopping: abandoned, not failed; nothing is written or recorded + std::lock_guard lock(m_Mutex); + m_Active--; + m_MemoryInUse -= std::min(m_MemoryInUse, job.memory); + continue; } catch (const std::exception& ex) { done.outcome.ok = false; done.outcome.error = std::string("crashed: ") + ex.what(); @@ -315,7 +329,11 @@ void UgcProcessor::Worker() { done.outcome.files.clear(); done.milliseconds = std::chrono::duration(std::chrono::steady_clock::now() - start).count(); done.cpuMilliseconds = std::max(0.0, UgcThrottle::ThreadCpuSeconds() - cpuStart) * 1000.0; - UgcThrottle::Checkpoint(); + try { + UgcThrottle::Checkpoint(); + } catch (const UgcThrottle::Cancelled&) { + // Stopping: the job is done already, so it's still recorded + } { std::lock_guard lock(m_Mutex); diff --git a/dUgcServer/UgcThrottle.cpp b/dUgcServer/UgcThrottle.cpp index 58c95acbd..4bd104860 100644 --- a/dUgcServer/UgcThrottle.cpp +++ b/dUgcServer/UgcThrottle.cpp @@ -14,6 +14,7 @@ namespace { constexpr double MIN_ACCOUNT_SECONDS = 0.005; std::atomic g_Budget{ 0.0 }; + std::atomic g_Cancel{ false }; std::mutex g_Mutex; double g_Balance = BURST_SECONDS; // CPU seconds that may still be used std::chrono::steady_clock::time_point g_Refilled = std::chrono::steady_clock::now(); @@ -58,7 +59,11 @@ namespace UgcThrottle { t_LastCpu = ThreadCpuSeconds(); } + void Cancel(const bool cancel) { g_Cancel = cancel; } + bool IsCancelled() { return g_Cancel; } + void Checkpoint() { + if (g_Cancel) throw Cancelled{}; const double budget = g_Budget; if (budget <= 0.0) return; const double cpu = ThreadCpuSeconds(); @@ -80,7 +85,12 @@ namespace UgcThrottle { wait = std::min(wait, 5.0); g_SleptMs += static_cast(wait * 1000.0); g_LastSleep = UnixMs(); - std::this_thread::sleep_for(std::chrono::duration(wait)); + // In short sleeps, so a cancel isn't held up by a long wait + const auto until = std::chrono::steady_clock::now() + std::chrono::duration_cast(std::chrono::duration(wait)); + while (std::chrono::steady_clock::now() < until) { + if (g_Cancel) throw Cancelled{}; + std::this_thread::sleep_for(std::min(until - std::chrono::steady_clock::now(), std::chrono::milliseconds(100))); + } // Time asleep costs no CPU; don't count this call's own bookkeeping twice t_LastCpu = ThreadCpuSeconds(); } diff --git a/dUgcServer/UgcThrottle.h b/dUgcServer/UgcThrottle.h index 634519149..4a658bc72 100644 --- a/dUgcServer/UgcThrottle.h +++ b/dUgcServer/UgcThrottle.h @@ -14,9 +14,16 @@ namespace UgcThrottle { void SetBudget(double cpus); double GetBudget(); - // Account the calling thread's CPU time and sleep when over the budget + // Account the calling thread's CPU time and sleep when over the budget. Throws Cancelled once Cancel(true) was + // called, so a job stops in moments when the server shuts down. void Checkpoint(); + // Thrown by Checkpoint while cancelled: the job is abandoned, not failed (its row stays waiting) + struct Cancelled {}; + // Stop (true) or allow (false) every job's work: set when the server stops, cleared when it starts + void Cancel(bool cancel); + bool IsCancelled(); + // Start accounting on this thread from now (a worker starting a job), so time spent idle isn't counted void Begin(); diff --git a/tests/dUgcTests/UgcTests.cpp b/tests/dUgcTests/UgcTests.cpp index 5df92a6e7..9ad6f2e0d 100644 --- a/tests/dUgcTests/UgcTests.cpp +++ b/tests/dUgcTests/UgcTests.cpp @@ -1755,3 +1755,12 @@ TEST(UgcHsr, SamplePointsFollowTheTrianglesSize) { EXPECT_LE(dense.size(), 40u); EXPECT_EQ(UgcHsr::SamplePoints({ 0, 0, 0 }, { 1.6f, 0, 0 }, { 1.6f, 0, 1.6f }, 0.1143f, 28).size(), 100u); // bigger ones as before } + +// Stopping the server cancels the jobs being made: Checkpoint throws until the cancel is cleared +TEST(UgcThrottle, CancelStopsJobsAtTheirNextCheckpoint) { + UgcThrottle::Cancel(true); + EXPECT_TRUE(UgcThrottle::IsCancelled()); + EXPECT_THROW(UgcThrottle::Checkpoint(), UgcThrottle::Cancelled); + UgcThrottle::Cancel(false); + EXPECT_NO_THROW(UgcThrottle::Checkpoint()); +}