fix(ugc): purging and the storage cap delete folders off the main thread

Purge all ran on the UGC server's main thread inside the dashboard's
request: it measured and deleted every stored folder, then updated rows one
by one, and the main thread (downloads, the master link, the database)
waited minutes for it. Now a purge answers at once and a file thread lists,
measures and deletes the folders; the main thread applies what it removed
and updates the rows a few hundred a tick. No new item is started while a
purge runs, and items being made when it starts are left alone, so no worker
writes into a folder being deleted. One purge at a time. The storage cap's
clean-up lists and deletes on the same thread. The UGC page shows the
purge's progress (/status purge).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Aaron Kimbrell
2026-09-28 19:19:07 -05:00
parent 3c004fd36d
commit bdb4c39db3
4 changed files with 245 additions and 84 deletions

View File

@@ -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' : '')) +
'</div>' + (l.nice ? '<div class="text-body-secondary">Workers run at priority ' + esc(l.nice) + '.</div>' : '') +
purgeLine(s.purge) +
(last ? '<div class="text-body-secondary">Last: ' + esc(last.kind) + ' ' + esc(last.id) + (last.ok ? ' made in ' + esc(duration(last.ms)) : ' failed') +
(last.message ? ' (' + esc(last.message) + ')' : '') + '</div>' : '');
}).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 '<div class="text-body-secondary">Last purge (' + what + '): ' + counts + ', finished ' + esc(new Date(p.finished * 1000).toLocaleString()) + '.</div>';
return '<div>' + fmt.badge(p.state === 'saving' ? 'Purge: updating the rows' : 'Purging', 'info') + ' ' + what + ': ' + counts + '. New items wait until it is done.</div>';
}
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();

View File

@@ -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<unsigned long long>(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::filesystem::file_time_type::duration>(std::chrono::hours(24 * count)); };
// What to delete: the files' ids, and the rows they came from (when they're known)
std::map<LWOOBJID, std::vector<LWOOBJID>> 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<LWOOBJID> 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<int64_t>(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::filesystem::file_time_type::duration>(std::chrono::hours(24 * count)); };
std::vector<LWOOBJID> 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> 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<unsigned long long>(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<int64_t>(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<unsigned long long>(m_PurgeInfo.deleted), static_cast<unsigned long long>(m_PurgeInfo.bytes));
}

View File

@@ -124,11 +124,22 @@ public:
uint64_t bytes{};
size_t busy{}; // being made right now, left alone
std::vector<std::string> 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<std::thread> 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<Removed> 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<LWOOBJID> m_PurgeSkip; // storage ids being made when it started
std::vector<std::pair<LWOOBJID, std::vector<LWOOBJID>>> m_PurgeTargets; // storage id -> rows (empty with `all`)
std::map<LWOOBJID, std::vector<LWOOBJID>> m_PurgeRowsOf; // main thread: rows of each target
std::vector<LWOOBJID> 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<std::pair<Kind, LWOOBJID>> m_InFlight; // models, and builds waiting for their combination
std::set<LWOOBJID> m_ComboJobs; // combinations being made

View File

@@ -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<unsigned long long>(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 = {},