diff --git a/dDashboardServer/static/js/ugc.js b/dDashboardServer/static/js/ugc.js index 4f4b93479..9c08271dc 100644 --- a/dDashboardServer/static/js/ugc.js +++ b/dDashboardServer/static/js/ugc.js @@ -314,6 +314,7 @@ meter('Jobs\' memory', s.jobMemoryBytes || 0, l.maxMemoryBytes || 0, esc(mb(s.jobMemoryBytes || 0)) + ' estimated' + (l.maxMemoryBytes ? ' of ' + esc(mb(l.maxMemoryBytes)) : ', no limit') + (s.memoryWaits ? ', ' + esc(s.memoryWaits) + ' waits' : '')) + '' + (l.nice ? '
Workers run at priority ' + esc(l.nice) + '.
' : '') + + purgeLine(s.purge) + (last ? '
Last: ' + esc(last.kind) + ' ' + esc(last.id) + (last.ok ? ' made in ' + esc(duration(last.ms)) : ' failed') + (last.message ? ' (' + esc(last.message) + ')' : '') + '
' : ''); }).catch(function () { box.innerHTML = fmt.badge('Unknown', 'secondary') + ' Couldn\'t ask the dashboard for the UGC server\'s status.'; }); @@ -597,9 +598,19 @@ // ---- deleting and purging stored files ---- + // The running or last purge (the UGC server deletes folders in the background) + function purgeLine(p) { + if (!p || !p.state) return ''; + var what = (p.kind === 'modular' ? 'cars and rockets' : 'models'); + var counts = esc(p.deleted) + ' deleted (' + esc(mb(p.bytes || 0)) + '), ' + esc(p.checked) + (p.total ? ' of ' + esc(p.total) : '') + ' looked at'; + if (p.state === 'done') return '
Last purge (' + what + '): ' + counts + ', finished ' + esc(new Date(p.finished * 1000).toLocaleString()) + '.
'; + return '
' + fmt.badge(p.state === 'saving' ? 'Purge: updating the rows' : 'Purging', 'info') + ' ' + what + ': ' + counts + '. New items wait until it is done.
'; + } + function showDeleteResult(box, d) { if (!d.success) { toast(d.error || 'Failed', 'danger'); return; } - var text = d.deleted + ' deleted (' + mb(d.bytes || 0) + ' freed). ' + (d.notes || []).join(' '); + var text = d.started ? 'Deleting ' + (d.queued ? d.queued + ' item(s)' : 'every stored item') + ' in the background. ' + (d.notes || []).join(' ') + : (d.notes && d.notes.length && !d.deleted ? d.notes.join(' ') : d.deleted + ' deleted (' + mb(d.bytes || 0) + ' freed). ' + (d.notes || []).join(' ')); if (box) box.textContent = text; toast(text, 'success'); load(); diff --git a/dUgcServer/UgcProcessor.cpp b/dUgcServer/UgcProcessor.cpp index 3048feba4..6481e9b42 100644 --- a/dUgcServer/UgcProcessor.cpp +++ b/dUgcServer/UgcProcessor.cpp @@ -147,7 +147,7 @@ void UgcProcessor::Stop() { for (auto& thread : m_Threads) { if (thread.joinable()) thread.join(); } - m_Threads.clear(); + m_Threads.clear(); if (m_FileThread.joinable()) m_FileThread.join(); } void UgcProcessor::Drain() { @@ -332,6 +332,8 @@ void UgcProcessor::Worker() { } void UgcProcessor::Poll() { + // No new item while a purge deletes folders (a worker would write into one being deleted) + if (m_FileTask == eFileTask::PURGE) return; size_t queued = 0, active = 0; Limits limits; UgcJobs::Settings settings; @@ -584,16 +586,12 @@ void UgcProcessor::Update() { m_NextPoll = now + std::chrono::milliseconds(m_Config.pollIntervalMs); Poll(); } - if (m_Config.maxStorageBytes > 0 && (m_StoredBytes > m_Config.maxStorageBytes || now >= m_NextEviction)) { + CollectFileTask(); + if (m_FileTask == eFileTask::NONE && m_Config.maxStorageBytes > 0 && (m_StoredBytes > m_Config.maxStorageBytes || now >= m_NextEviction)) { m_NextEviction = now + EVICTION_INTERVAL; - // Deleted files stay marked made: they're made again when someone asks for them (Request) - const auto removed = m_Storage.Evict(m_Config.maxStorageBytes); - m_StoredBytes = 0; - for (const auto& entry : m_Storage.List()) m_StoredBytes += entry.bytes; - if (!removed.empty()) { - m_Evicted += removed.size(); - LOG("Deleted the files of %zu item(s) used longest ago to stay under %llu MB", removed.size(), static_cast(m_Config.maxStorageBytes / (1024 * 1024))); - } + // Deleted files stay marked made: they're made again when someone asks for them (Request). The folders are + // listed and deleted on the file thread. + StartFileTask(eFileTask::EVICT); } // Forget old answers std::erase_if(m_Recent, [now](const auto& item) { return now - item.second.first > RECENT_ANSWER_TIME; }); @@ -642,6 +640,7 @@ nlohmann::json UgcProcessor::Status() const { { std::lock_guard lock(m_Mutex); status["queued"] = m_Jobs.size(); + status["purge"] = PurgeStatus(); status["active"] = m_Active; } Limits limits; @@ -815,80 +814,184 @@ size_t UgcProcessor::RegenerateIcons(const std::string& kind) { UgcProcessor::DeleteResult UgcProcessor::Delete(const DeleteRequest& request) { DeleteResult result; - const auto now = std::filesystem::file_time_type::clock::now(); - const auto days = [](int64_t count) { return std::chrono::duration_cast(std::chrono::hours(24 * count)); }; - - // What to delete: the files' ids, and the rows they came from (when they're known) - std::map> targets; // storage id -> rows - if (request.all) { - for (const auto& entry : m_Storage.List()) { - if (entry.kind == request.kind) targets[entry.id]; - } + if (m_FileTask != eFileTask::NONE) { + result.notes.push_back(m_FileTask == eFileTask::PURGE ? "A purge is already running: try again when it's done (the UGC page shows it)." : + "The storage cap's clean-up is running: try again in a moment."); + return result; } + m_Purge = request; + m_PurgeSkip.clear(); + for (const auto& [kind, id] : m_InFlight) if (kind == request.kind) m_PurgeSkip.insert(id); + if (request.kind == Kind::MODULAR) m_PurgeSkip.insert(m_ComboJobs.begin(), m_ComboJobs.end()); + m_PurgeTargets.clear(); + m_PurgeRowsOf.clear(); for (const auto id : request.ids) { const auto storageId = StorageId(request.kind, id); - if (storageId != 0) targets[storageId].push_back(id); + if (storageId != 0) m_PurgeRowsOf[storageId].push_back(id); } - - std::vector deletedRows; - for (const auto& [storageId, rows] : targets) { - const bool busy = request.kind == Kind::MODEL ? m_InFlight.contains({ Kind::MODEL, storageId }) : m_ComboJobs.contains(storageId); - if (busy) { - result.busy++; - continue; - } - const auto folder = m_Storage.Folder(request.kind, storageId); - std::error_code error; - if (!std::filesystem::exists(folder, error)) { - deletedRows.insert(deletedRows.end(), rows.begin(), rows.end()); - continue; - } - if (request.unusedDays > 0) { - const auto used = std::filesystem::last_write_time(folder, error); - if (!error && now - used < days(request.unusedDays)) continue; - } - if (request.olderThanDays > 0) { - const auto made = std::filesystem::last_write_time(folder / "icon.png", error); - if (!error && now - made < days(request.olderThanDays)) continue; - } - uint64_t bytes = 0; - for (const auto& file : std::filesystem::directory_iterator(folder, error)) bytes += file.is_regular_file(error) ? file.file_size(error) : 0; - m_Storage.Remove(request.kind, storageId); - result.deleted++; - result.bytes += bytes; - m_StoredBytes -= std::min(m_StoredBytes, bytes); - m_Recent.clear(); - deletedRows.insert(deletedRows.end(), rows.begin(), rows.end()); - if (request.kind == Kind::MODEL) deletedRows.push_back(storageId); + for (const auto& [storageId, rows] : m_PurgeRowsOf) m_PurgeTargets.emplace_back(storageId, rows); + if (!request.all && m_PurgeTargets.empty()) { + result.notes.push_back("Nothing to delete."); + return result; } - std::sort(deletedRows.begin(), deletedRows.end()); - deletedRows.erase(std::unique(deletedRows.begin(), deletedRows.end()), deletedRows.end()); - - // What happens to the rows - if (request.after == eAfterDelete::NOW) { - if (request.all) { - if (request.kind == Kind::MODEL) Database::Get()->ResetUgcModelProcessing(std::nullopt, false); - else Database::Get()->ResetModularBuildProcessing(std::nullopt, false); - } else { - for (const auto row : deletedRows) { - if (request.kind == Kind::MODEL) Database::Get()->ResetUgcModelProcessing(row, false); - else Database::Get()->ResetModularBuildProcessing(row, false); - } - } - m_NextPoll = std::chrono::steady_clock::now(); - result.notes.push_back("They're queued to be made again now."); - } else if (request.after == eAfterDelete::GONE) { - for (const auto row : deletedRows) { - if (request.kind == Kind::MODEL) Database::Get()->SetUgcModelProcessed(row, IUgc::eProcessState::FAILED, m_Config.maxAttempts, "Deleted from the dashboard", false); - else Database::Get()->SetModularBuildProcessed(row, IUgc::eProcessState::FAILED, m_Config.maxAttempts, "Deleted from the dashboard"); - } - result.notes.push_back("They're marked deleted and won't be made again unless made again from the dashboard."); - } else { - result.notes.push_back("They'll be made again when a game client asks for them."); - } - if (request.kind == Kind::MODULAR && result.deleted > 0) { - result.notes.push_back("Car and rocket icons are shared by every build of the same modules: the other builds' icons are made again when asked for."); - } - if (result.busy > 0) result.notes.push_back(std::to_string(result.busy) + " being made right now were left alone."); + m_PurgeRows.clear(); + m_PurgeRowsDone = 0; + m_PurgeInfo = PurgeInfo{ "running", request.kind == Kind::MODULAR, 0, request.all ? 0 : m_PurgeTargets.size(), 0, 0, static_cast(std::time(nullptr)), 0 }; + StartFileTask(eFileTask::PURGE); + result.started = true; + result.queued = request.all ? 0 : m_PurgeTargets.size(); + result.notes.push_back("Deleting in the background; no new item is made until it's done. The UGC page shows how far it got."); return result; } + +nlohmann::json UgcProcessor::PurgeStatus() const { + if (m_PurgeInfo.state.empty()) return nlohmann::json::object(); + return { { "state", m_PurgeInfo.state }, { "kind", m_PurgeInfo.modular ? "modular" : "model" }, { "checked", m_PurgeInfo.checked }, { "total", m_PurgeInfo.total }, + { "deleted", m_PurgeInfo.deleted }, { "bytes", m_PurgeInfo.bytes }, { "started", m_PurgeInfo.started }, { "finished", m_PurgeInfo.finished } }; +} + +void UgcProcessor::StartFileTask(const eFileTask type) { + { + std::lock_guard lock(m_FileMutex); + m_FileRemoved.clear(); + m_FileChecked = 0; + m_FileStoredBytes = 0; + m_FileDone = false; + } + m_FileTask = type; + if (m_FileThread.joinable()) m_FileThread.join(); + m_FileThread = std::thread([this, type, request = m_Purge, skip = m_PurgeSkip, targets = m_PurgeTargets, maxBytes = m_Config.maxStorageBytes]() { + const auto remove = [this](LWOOBJID storageId, uint64_t bytes, Kind kind) { + m_Storage.Remove(kind, storageId); + std::lock_guard lock(m_FileMutex); + m_FileRemoved.push_back({ storageId, bytes }); + }; + if (type == eFileTask::EVICT) { + // The storage cap: the items used longest ago go until the rest fits (UgcStorage::Evict, on this thread) + auto entries = m_Storage.List(); + uint64_t total = 0; + for (const auto& entry : entries) total += entry.bytes; + if (total > maxBytes) { + std::sort(entries.begin(), entries.end(), [](const auto& a, const auto& b) { return a.used < b.used; }); + for (const auto& entry : entries) { + if (total <= maxBytes) break; + remove(entry.id, entry.bytes, entry.kind); + total -= std::min(total, entry.bytes); + } + } + std::lock_guard lock(m_FileMutex); + m_FileStoredBytes = total; + m_FileDone = true; + return; + } + const auto now = std::filesystem::file_time_type::clock::now(); + const auto days = [](int64_t count) { return std::chrono::duration_cast(std::chrono::hours(24 * count)); }; + std::vector ids; + if (request.all) { + for (const auto& entry : m_Storage.List()) if (entry.kind == request.kind) ids.push_back(entry.id); + } else { + for (const auto& [storageId, rows] : targets) ids.push_back(storageId); + } + for (const auto storageId : ids) { + { + std::lock_guard lock(m_FileMutex); + m_FileChecked++; + } + if (skip.contains(storageId)) continue; + const auto folder = m_Storage.Folder(request.kind, storageId); + std::error_code error; + if (!std::filesystem::exists(folder, error)) { + // Nothing stored: its rows still get what was asked for + std::lock_guard lock(m_FileMutex); + m_FileRemoved.push_back({ storageId, 0 }); + continue; + } + if (request.unusedDays > 0) { + const auto used = std::filesystem::last_write_time(folder, error); + if (!error && now - used < days(request.unusedDays)) continue; + } + if (request.olderThanDays > 0) { + const auto made = std::filesystem::last_write_time(folder / "icon.png", error); + if (!error && now - made < days(request.olderThanDays)) continue; + } + uint64_t bytes = 0; + for (const auto& file : std::filesystem::directory_iterator(folder, error)) bytes += file.is_regular_file(error) ? file.file_size(error) : 0; + remove(storageId, bytes, request.kind); + } + std::lock_guard lock(m_FileMutex); + m_FileDone = true; + }); +} + +void UgcProcessor::CollectFileTask() { + if (m_FileTask == eFileTask::NONE) return; + std::vector removed; + bool done = false; + size_t checked = 0; + uint64_t storedLeft = 0; + { + std::lock_guard lock(m_FileMutex); + removed.swap(m_FileRemoved); + done = m_FileDone; + checked = m_FileChecked; + storedLeft = m_FileStoredBytes; + } + if (!removed.empty()) m_Recent.clear(); + if (m_FileTask == eFileTask::EVICT) { + for (const auto& item : removed) m_StoredBytes -= std::min(m_StoredBytes, item.bytes); + m_Evicted += removed.size(); + m_EvictedThisTask += removed.size(); + if (!done) return; + m_FileThread.join(); + m_StoredBytes = storedLeft; + m_FileTask = eFileTask::NONE; + if (m_EvictedThisTask > 0) { + LOG("Deleted the files of %zu item(s) used longest ago to stay under %llu MB", m_EvictedThisTask, static_cast(m_Config.maxStorageBytes / (1024 * 1024))); + } + m_EvictedThisTask = 0; + return; + } + + // A purge: what the thread removed so far, then the rows (a few hundred a tick) + const auto kind = m_Purge.kind; + for (const auto& item : removed) { + m_StoredBytes -= std::min(m_StoredBytes, item.bytes); + if (item.bytes > 0) { + m_PurgeInfo.deleted++; + m_PurgeInfo.bytes += item.bytes; + } + if (const auto it = m_PurgeRowsOf.find(item.storageId); it != m_PurgeRowsOf.end()) m_PurgeRows.insert(m_PurgeRows.end(), it->second.begin(), it->second.end()); + if (kind == Kind::MODEL) m_PurgeRows.push_back(item.storageId); + } + m_PurgeInfo.checked = checked; + if (!done) return; + if (m_FileThread.joinable()) m_FileThread.join(); + m_PurgeInfo.state = "saving"; + + constexpr size_t ROWS_PER_TICK = 250; + if (m_Purge.after == eAfterDelete::NOW && m_Purge.all) { + if (kind == Kind::MODEL) Database::Get()->ResetUgcModelProcessing(std::nullopt, false); + else Database::Get()->ResetModularBuildProcessing(std::nullopt, false); + m_PurgeRowsDone = m_PurgeRows.size(); + } else if (m_Purge.after != eAfterDelete::ON_DEMAND) { + std::sort(m_PurgeRows.begin(), m_PurgeRows.end()); + m_PurgeRows.erase(std::unique(m_PurgeRows.begin(), m_PurgeRows.end()), m_PurgeRows.end()); + const auto end = std::min(m_PurgeRows.size(), m_PurgeRowsDone + ROWS_PER_TICK); + for (; m_PurgeRowsDone < end; m_PurgeRowsDone++) { + const auto row = m_PurgeRows[m_PurgeRowsDone]; + if (m_Purge.after == eAfterDelete::NOW) { + if (kind == Kind::MODEL) Database::Get()->ResetUgcModelProcessing(row, false); + else Database::Get()->ResetModularBuildProcessing(row, false); + } else { + if (kind == Kind::MODEL) Database::Get()->SetUgcModelProcessed(row, IUgc::eProcessState::FAILED, m_Config.maxAttempts, "Deleted from the dashboard", false); + else Database::Get()->SetModularBuildProcessed(row, IUgc::eProcessState::FAILED, m_Config.maxAttempts, "Deleted from the dashboard"); + } + } + if (m_PurgeRowsDone < m_PurgeRows.size()) return; + } + m_PurgeInfo.state = "done"; + m_PurgeInfo.finished = static_cast(std::time(nullptr)); + m_FileTask = eFileTask::NONE; + if (m_Purge.after == eAfterDelete::NOW) m_NextPoll = std::chrono::steady_clock::now(); + LOG("Purge done: the files of %llu item(s) deleted (%llu bytes)", static_cast(m_PurgeInfo.deleted), static_cast(m_PurgeInfo.bytes)); +} diff --git a/dUgcServer/UgcProcessor.h b/dUgcServer/UgcProcessor.h index ac8a143a3..9ec5be119 100644 --- a/dUgcServer/UgcProcessor.h +++ b/dUgcServer/UgcProcessor.h @@ -124,11 +124,22 @@ public: uint64_t bytes{}; size_t busy{}; // being made right now, left alone std::vector notes; + bool started{}; // the files are being deleted in the background (PurgeStatus); `deleted` and `bytes` come later + size_t queued{}; // items to look at (0 with `all`: every stored item of the kind) }; - // Main thread: deletes stored files. Items being made are skipped, so a worker never writes what is being deleted. + /** + * Main thread: deletes stored files, on a background thread (a purge can be thousands of folders, and the main + * thread serves downloads and the master link): the reply comes at once, and PurgeStatus follows it. Items being + * made are skipped, and no new item is started until the purge is done, so a worker never writes what is being + * deleted. One purge (or storage clean-up) at a time; another is refused while one runs. + */ DeleteResult Delete(const DeleteRequest& request); + // Main thread: the running or last purge ({} when none ran): {state: running|saving|done, kind, checked, deleted, + // bytes, started (Unix seconds), finished} + nlohmann::json PurgeStatus() const; + // Main thread: what the server is doing, for /status nlohmann::json Status() const; @@ -231,6 +242,41 @@ private: bool m_Draining{}; // main thread: live update, no new work (Drain) std::vector m_Threads; + // Deleting stored folders off the main thread: a purge (Delete) or the storage cap (Evict). The thread only touches + // the filesystem; the main thread applies what it removed (Update -> CollectFileTask). + enum class eFileTask { NONE, PURGE, EVICT }; + struct Removed { + LWOOBJID storageId{}; + uint64_t bytes{}; + }; + void StartFileTask(eFileTask type); + void CollectFileTask(); + eFileTask m_FileTask{ eFileTask::NONE }; // main thread + std::thread m_FileThread; + mutable std::mutex m_FileMutex; // guards the members down to m_FileDone + std::vector m_FileRemoved; // removed since the main thread last looked + size_t m_FileChecked{}; + uint64_t m_FileStoredBytes{}; // EVICT: the total left + bool m_FileDone{}; + DeleteRequest m_Purge; // the purge running (main thread, read by the thread before it starts) + std::set m_PurgeSkip; // storage ids being made when it started + std::vector>> m_PurgeTargets; // storage id -> rows (empty with `all`) + std::map> m_PurgeRowsOf; // main thread: rows of each target + std::vector m_PurgeRows; // main thread: rows whose files are gone, to update + size_t m_PurgeRowsDone{}; + struct PurgeInfo { + std::string state; // empty: none ran; running, saving (the rows), done + bool modular{}; + size_t checked{}; + size_t total{}; // 0 with `all` + uint64_t deleted{}; + uint64_t bytes{}; + int64_t started{}; + int64_t finished{}; + }; + PurgeInfo m_PurgeInfo; // main thread: PurgeStatus + size_t m_EvictedThisTask{}; // main thread + // Main thread only std::set> m_InFlight; // models, and builds waiting for their combination std::set m_ComboJobs; // combinations being made diff --git a/dUgcServer/UgcServer.cpp b/dUgcServer/UgcServer.cpp index adaee0abf..6cc43de8a 100644 --- a/dUgcServer/UgcServer.cpp +++ b/dUgcServer/UgcServer.cpp @@ -495,9 +495,10 @@ namespace { const auto after = body->value("after", std::string("on_demand")); request.after = after == "now" ? UgcProcessor::eAfterDelete::NOW : after == "gone" ? UgcProcessor::eAfterDelete::GONE : UgcProcessor::eAfterDelete::ON_DEMAND; const auto result = g_Processor->Delete(request); - LOG("Deleted the files of %zu item(s) (%llu bytes) for the dashboard", result.deleted, static_cast(result.bytes)); + LOG(result.started ? "Purge started for the dashboard (%zu item(s) to look at, 0: all)" : "Purge not started: %zu", result.started ? result.queued : result.notes.size()); reply.status = eHTTPStatusCode::OK; - reply.message = nlohmann::json{ { "success", true }, { "deleted", result.deleted }, { "bytes", result.bytes }, { "busy", result.busy }, { "notes", result.notes } }.dump(); + reply.message = nlohmann::json{ { "success", true }, { "deleted", result.deleted }, { "bytes", result.bytes }, { "busy", result.busy }, { "notes", result.notes }, + { "started", result.started }, { "queued", result.queued } }.dump(); } }); Game::web.RegisterHTTPRoute({ .path = "/status", .method = eHTTPMethod::GET, .middleware = {},