#include "UgcProcessor.h" #include "CDClientDatabase.h" #include "Database.h" #include "Logger.h" #include "UgcBricks.h" #include "UgcCdClient.h" #include "Sd0.h" #include "UgcFormats.h" #include "UgcKeys.h" #include "ZCompression.h" #include "UgcThrottle.h" #include "json.hpp" #include #include #include #include #include #if defined(__linux__) #include #include #include #include #endif namespace { constexpr size_t MAX_ERROR_LENGTH = 1000; // process_error holds 1024 constexpr size_t LOG_LENGTH = 50; constexpr auto EVICTION_INTERVAL = std::chrono::minutes(5); constexpr auto RECENT_ANSWER_TIME = std::chrono::seconds(10); // In the storage folder once every item stored has its sd0 icon and checksums (Backfill) constexpr auto BACKFILL_MARKER = ".checksums-stored"; constexpr auto STATS_BACKFILL_MARKER = ".stats-stored"; // What a model's stats.json counted: bricks, and the most detailed level's triangles after and before the faces that // can't be seen were removed (the transparent ones are kept as they are) struct ModelCounts { uint32_t bricks{}; uint32_t triangles{}; uint32_t trianglesBefore{}; }; std::optional CountsOf(const nlohmann::json& stats) { if (!stats.is_object()) return std::nullopt; const auto& lods = stats.value("lods", nlohmann::json::array()); const auto& first = lods.is_array() && !lods.empty() ? lods.front() : nlohmann::json::object(); const auto transparent = first.value("transparent", 0u); return ModelCounts{ stats.value("bricks", 0u), first.value("opaqueAfter", 0u) + transparent, first.value("opaqueBefore", 0u) + transparent }; } constexpr size_t BACKFILL_ITEMS_PER_UPDATE = 16; constexpr uint32_t BACKFILL_BUILDS_PER_UPDATE = 200; // Combination id of a build whose modules can't be told (so it isn't looked at again) constexpr LWOOBJID NO_COMBINATION = -1; const char* KindName(UgcStorage::Kind kind) { return kind == UgcStorage::Kind::MODEL ? "model" : "modular"; } // The process's CPU time (user and system, every thread), seconds double ProcessCpuSeconds() { timespec ts{}; if (clock_gettime(CLOCK_PROCESS_CPUTIME_ID, &ts) == 0) return static_cast(ts.tv_sec) + ts.tv_nsec / 1e9; return 0.0; } void ApplyNice(int nice) { #if defined(__linux__) // Per thread on Linux: only the calling worker's priority changes setpriority(PRIO_PROCESS, static_cast(syscall(SYS_gettid)), std::clamp(nice, 0, 19)); #else (void)nice; #endif } int64_t UnixNow() { return std::chrono::duration_cast(std::chrono::system_clock::now().time_since_epoch()).count(); } } UgcProcessor::UgcProcessor(Config config, UgcStorage& storage, UgcBricks::BrickLibrary& library, UgcJobs::Settings settings) : m_Config(config), m_Storage(storage), m_Library(library), m_Settings(std::move(settings)) {} uint64_t UgcProcessor::ResidentBytes() { #if defined(__linux__) std::ifstream statm("/proc/self/statm"); uint64_t size = 0, resident = 0; if (statm >> size >> resident) return resident * static_cast(sysconf(_SC_PAGESIZE)); #endif return 0; } bool UgcProcessor::Throttled() const { const auto last = UgcThrottle::GetStats().lastSleepUnixMs; return last > 0 && UnixNow() * 1000 - last < 5000; } void UgcProcessor::Configure(UgcJobs::Settings settings, Limits limits) { { std::lock_guard lock(m_Mutex); m_Settings = std::move(settings); m_Limits = limits; } UgcThrottle::SetBudget(limits.maxCpus); m_Wake.notify_all(); } void UgcProcessor::SampleUsage() { const auto now = std::chrono::steady_clock::now(); const double cpu = ProcessCpuSeconds(); if (m_CpuSampled.time_since_epoch().count() != 0) { const double elapsed = std::chrono::duration(now - m_CpuSampled).count(); if (elapsed > 0.0) m_CpuPercent = (cpu - m_CpuSeconds) / elapsed * 100.0; } m_CpuSampled = now; m_CpuSeconds = cpu; } UgcProcessor::~UgcProcessor() { Stop(); } void UgcProcessor::Start() { const auto stored = m_Storage.List(); for (const auto& entry : stored) m_StoredBytes += entry.bytes; // Items made before the checksums were stored get their sd0 icon and checksums once (Backfill) std::error_code error; m_BackfillItemsDone = std::filesystem::exists(m_Storage.GetRoot() / BACKFILL_MARKER, error); if (!m_BackfillItemsDone) m_BackfillItems.assign(stored.begin(), stored.end()); // Models made before the triangles before hidden-face removal were stored get them from their stats.json once m_StatsBackfillDone = std::filesystem::exists(m_Storage.GetRoot() / STATS_BACKFILL_MARKER, error); if (!m_StatsBackfillDone) { for (const auto& entry : stored) if (entry.kind == Kind::MODEL) m_StatsBackfill.push_back(entry.id); } 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))); } void UgcProcessor::Stop() { { std::lock_guard lock(m_Mutex); 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(); } m_Threads.clear(); if (m_FileThread.joinable()) m_FileThread.join(); } void UgcProcessor::Drain() { if (m_Draining) return; m_Draining = true; std::deque dropped; { std::lock_guard lock(m_Mutex); dropped.swap(m_Jobs); } // What was waiting for them is forgotten too; the database rows stay pending for (const auto& job : dropped) { if (job.kind == Kind::MODEL) { m_InFlight.erase({ job.kind, job.id }); continue; } m_ComboJobs.erase(job.id); if (const auto rows = m_ComboRows.find(job.id); rows != m_ComboRows.end()) { for (const auto& [row, attempts] : rows->second) m_InFlight.erase({ Kind::MODULAR, row }); m_ComboRows.erase(rows); } } LOG("Draining: %zu queued job(s) left for the next UGC server, waiting for the running ones", dropped.size()); } bool UgcProcessor::Drained() const { if (!m_Draining) return false; std::lock_guard lock(m_Mutex); return m_Active == 0 && m_Jobs.empty() && m_Done.empty(); } void UgcProcessor::Worker() { int appliedNice = 0; while (true) { Job job; UgcJobs::Settings settings; int nice = 0; { std::unique_lock lock(m_Mutex); // The first job that fits the memory budget beside the ones running; one that is too big alone runs // when nothing else does auto pick = m_Jobs.end(); m_Wake.wait(lock, [this, &pick] { if (m_Stopping) return true; pick = m_Jobs.end(); for (auto it = m_Jobs.begin(); it != m_Jobs.end(); ++it) { if (m_Limits.maxMemoryBytes == 0 || m_Active == 0 || m_MemoryInUse + it->memory <= m_Limits.maxMemoryBytes) { pick = it; break; } } if (pick == m_Jobs.end() && !m_Jobs.empty()) m_Waiting++; return pick != m_Jobs.end(); }); if (m_Stopping) return; job = std::move(*pick); m_Jobs.erase(pick); m_Active++; m_MemoryInUse += job.memory; settings = m_Settings; nice = m_Limits.nice; } if (nice != appliedNice) { ApplyNice(nice); appliedNice = nice; } UgcThrottle::Begin(); if (job.preview) { // A preview: the icon goes back to whoever asked, nothing is stored or recorded HTTPReply reply; try { UgcJobs::Outcome outcome{ false, "cancelled" }; if (job.assembly) { if (!job.preview.Cancelled()) { auto nif = UgcJobs::AssemblyNif(job.modular, m_Library.GetResPath(), outcome.error); if (nif) { auto shared = std::make_shared(std::move(*nif)); CacheAssembly(job.modular.key, shared); outcome.ok = true; outcome.files["assembly.nif"] = *shared; } } } else if (!job.preview.Cancelled() && job.kind == Kind::MODULAR) { outcome = UgcJobs::ProcessModular(job.modular, m_Library.GetResPath(), settings); } else if (!job.preview.Cancelled()) { // A player model's icon from its stored .nif const auto nif = m_Storage.ReadNif(Kind::MODEL, job.id, "model.nif"); auto options = settings.icon; UgcIconParams::Apply(options, job.iconValues); const auto plain = options.denoise != UgcRender::eDenoise::OFF ? m_Storage.ReadNif(Kind::MODEL, job.id, "model.noao.nif") : std::nullopt; outcome.ok = nif && UgcJobs::IconFromNif(*nif, options, outcome.files, outcome.error, settings.shaders.TagLooks(), settings.shaders.OverlayTags(), plain ? &*plain : nullptr); if (!nif) outcome.error = "the model has no stored .nif yet"; } if (outcome.ok && outcome.files.contains("assembly.nif")) { reply.status = eHTTPStatusCode::OK; reply.contentType = eContentType::APPLICATION_OCTET_STREAM; reply.message = std::move(outcome.files["assembly.nif"]); reply.headers.push_back("Cache-Control: no-store"); } else if (outcome.ok && outcome.files.contains("icon.png")) { reply.status = eHTTPStatusCode::OK; reply.contentType = eContentType::IMAGE_PNG; reply.message = std::move(outcome.files["icon.png"]); reply.headers.push_back("Cache-Control: no-store"); } else { reply.status = eHTTPStatusCode::UNPROCESSABLE_ENTITY; 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; reply.message = ex.what(); } job.preview.Send(std::move(reply)); { std::lock_guard lock(m_Mutex); m_Active--; m_MemoryInUse -= std::min(m_MemoryInUse, job.memory); } m_Wake.notify_all(); continue; } const auto start = std::chrono::steady_clock::now(); const double cpuStart = UgcThrottle::JobCpuSeconds(); Done done{ job.kind, job.id, job.attempts }; done.memoryEstimate = job.memory; done.iconOnly = job.iconOnly; try { if (job.iconOnly) { // Only the icon, from the .nif made before const auto nif = m_Storage.ReadNif(Kind::MODEL, job.id, "model.nif"); auto options = settings.icon; UgcIconParams::Apply(options, job.iconValues); const auto iconStart = std::chrono::steady_clock::now(); const auto plain = options.denoise != UgcRender::eDenoise::OFF ? m_Storage.ReadNif(Kind::MODEL, job.id, "model.noao.nif") : std::nullopt; done.outcome.ok = nif && UgcJobs::IconFromNif(*nif, options, done.outcome.files, done.outcome.error, settings.shaders.TagLooks(), settings.shaders.OverlayTags(), plain ? &*plain : nullptr); if (!nif) done.outcome.error = "no stored .nif"; // The make's time keeps its icon's: the stats get the new icon's time, the row the difference (Collect) const auto stats = done.outcome.ok ? m_Storage.ReadNif(Kind::MODEL, job.id, "stats.json") : std::nullopt; const double iconMs = std::chrono::duration(std::chrono::steady_clock::now() - iconStart).count(); if (const auto updated = stats ? UgcJobs::WithIconTime(*stats, iconMs, done.iconChangeMs) : std::nullopt) done.outcome.files["stats.json"] = *updated; } else { // The processing options staff picked for this make, over the settings (unknown ones: the settings') UgcProcessOptions::Choice choice; if (job.kind == Kind::MODEL && UgcProcessOptions::Parse(job.options, choice)) UgcJobs::ApplyOptions(settings, choice); done.outcome = job.kind == Kind::MODEL ? 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(); } job.blob.clear(); job.blob.shrink_to_fit(); if (done.outcome.ok) { std::string error; const auto bytes = job.iconOnly ? m_Storage.Update(job.kind, job.id, done.outcome.files, error) : m_Storage.Write(job.kind, job.id, done.outcome.files, error); if (bytes) { done.bytes = *bytes; // The checksums of what the client downloads as sd0, for the main thread to store (the worlds answer // the clients' manifest requests with them) for (const std::string name : { "icon.dds", "model.nif" }) { const auto checksum = done.outcome.files.find(name + ".checksum"); if (checksum == done.outcome.files.end() || !done.outcome.files.contains(name + ".sd0")) continue; Checksum parsed{ name }; if (UgcFormats::ReadChecksumXml(checksum->second, parsed.md5, parsed.size)) done.checksums.push_back(std::move(parsed)); } } else { done.outcome.ok = false; done.outcome.error = error; } } done.outcome.files.clear(); done.milliseconds = std::chrono::duration(std::chrono::steady_clock::now() - start).count(); done.cpuMilliseconds = std::max(0.0, UgcThrottle::JobCpuSeconds() - cpuStart) * 1000.0; try { UgcThrottle::Checkpoint(); } catch (const UgcThrottle::Cancelled&) { // Stopping: the job is done already, so it's still recorded } { std::lock_guard lock(m_Mutex); m_Active--; m_MemoryInUse -= std::min(m_MemoryInUse, job.memory); m_Done.push_back(std::move(done)); } #if defined(__linux__) && defined(__GLIBC__) // Give the big buffers of the job back to the system malloc_trim(0); #endif m_Wake.notify_all(); } } 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; { std::lock_guard lock(m_Mutex); queued = m_Jobs.size(); active = m_Active; limits = m_Limits; settings = m_Settings; } const auto now = std::time(nullptr); std::tm local{}; #if defined(_WIN32) localtime_s(&local, &now); #else localtime_r(&now, &local); #endif m_Paused = UgcThrottle::InHours(local.tm_hour, limits.pauseFromHour, limits.pauseToHour); if (m_Paused) return; // Enough to keep every worker busy until the next poll const size_t wanted = std::max(m_Threads.size() * 2, m_Config.pollBatch); // Cars and rockets come first: the game client asks for their icons as soon as they're built. So a few are polled // even when the queue is full of models, and they go to the front of the queue. const bool full = queued + active >= wanted; const auto limit = full ? 0u : static_cast(wanted - queued - active + m_InFlight.size()); const auto buildLimit = std::max(limit, static_cast(m_Threads.size() + m_InFlight.size())); std::vector jobs; std::vector models; // Models staff asked to be made again come first too: polled even when the queue is full, to its front if (limit > 0) models = Database::Get()->GetUgcModelsToProcess(limit); else models = Database::Get()->GetUgcModelsToProcess(buildLimit, true); for (auto& model : models) { if (m_InFlight.contains({ Kind::MODEL, model.id })) continue; Job job{ Kind::MODEL, model.id, model.attempts, UgcJobs::LxfmlFromBlob(model.lxfml) }; job.priority = model.priority; job.options = std::move(model.options); if (job.blob.empty()) job.blob = std::move(model.lxfml); // the worker reports it can't be read job.parts = UgcJobs::CountParts(job.blob); job.iconValues = IconValues(UgcIconParams::ModelKind(), UgcIconParams::ModelTarget(model.id)); job.memory = UgcJobs::EstimateMemory(job.parts, settings); jobs.push_back(std::move(job)); } // Cars and rockets: one icon per combination of modules, shared by every build of it for (auto& build : Database::Get()->GetModularBuildsToProcess(buildLimit)) { if (m_InFlight.contains({ Kind::MODULAR, build.id })) continue; const auto key = UgcModularKey::Normalize(build.modules); if (key.empty()) { Record(Done{ Kind::MODULAR, build.id, build.attempts, UgcJobs::Outcome{ false, "no modules in \"" + build.modules + "\"" } }); continue; } const auto combo = UgcModularKey::StorageId(key); m_ComboOf[build.id] = combo; // Made already for another build of the same modules if (!m_ComboJobs.contains(combo) && m_Storage.File(Kind::MODULAR, combo, "icon.png")) { UgcJobs::Outcome reused{ true }; reused.note = "the same modules as a build made before (" + key + ")"; m_Reused++; Record(Done{ Kind::MODULAR, build.id, build.attempts, std::move(reused) }); continue; } m_InFlight.insert({ Kind::MODULAR, build.id }); m_ComboRows[combo].emplace_back(build.id, build.attempts); if (m_ComboJobs.contains(combo)) continue; Job job{ Kind::MODULAR, combo, build.attempts }; job.memory = UgcJobs::EstimateMemory(64, settings); std::string error; if (!UgcCdClient::GatherModular(build.modules, job.modular, error)) { // Nothing a worker could do: record it right away m_InFlight.erase({ Kind::MODULAR, build.id }); m_ComboRows.erase(combo); Record(Done{ Kind::MODULAR, build.id, build.attempts, UgcJobs::Outcome{ false, error } }); continue; } job.modular.key = key; job.modular.iconValues = IconValues(UgcIconParams::BuildKind(job.modular.buildType), UgcIconParams::CombinationTarget(key)); m_ComboJobs.insert(combo); jobs.push_back(std::move(job)); } if (jobs.empty()) return; { std::lock_guard lock(m_Mutex); for (auto& job : jobs) { if (job.kind == Kind::MODEL) m_InFlight.insert({ job.kind, job.id }); if (job.kind == Kind::MODULAR || job.priority) m_Jobs.push_front(std::move(job)); else m_Jobs.push_back(std::move(job)); } } m_Wake.notify_all(); } void UgcProcessor::RecordIconTime(LWOOBJID id, double changeMs) { for (const auto& entry : Database::Get()->GetUgcEntries({ id })) { if (entry.kind != IUgcLookup::eUgcKind::MODEL || entry.processMs == 0) continue; // The icon is drawn on one thread, so its time is its CPU time too const auto change = [changeMs](uint32_t value) { return static_cast(std::max(0.0, static_cast(value) + changeMs)); }; Database::Get()->SetUgcModelProcessStats(id, { change(entry.processMs), change(entry.processCpuMs), entry.processMemoryKb }); } } void UgcProcessor::Record(const Done& done) { // What the make cost (wall time, the worker's CPU time, the estimated memory), for the dashboard const IUgc::ProcessStats cost{ static_cast(done.milliseconds), static_cast(done.cpuMilliseconds), static_cast(done.memoryEstimate / 1024) }; auto error = done.outcome.error.substr(0, MAX_ERROR_LENGTH); const auto attempts = done.attempts + 1; const auto state = done.outcome.ok ? IUgc::eProcessState::DONE : done.outcome.empty ? IUgc::eProcessState::EMPTY : attempts >= m_Config.maxAttempts ? IUgc::eProcessState::FAILED : IUgc::eProcessState::PENDING; if (done.kind == Kind::MODEL) { Database::Get()->SetUgcModelProcessed(done.id, state, attempts, error, done.outcome.ok && done.outcome.aoBaked); if (done.outcome.ok) Database::Get()->SetUgcModelProcessStats(done.id, cost); // What it counted (stats.json), for sorting on the dashboard const auto stats = done.outcome.ok && !done.outcome.stats.empty() ? nlohmann::json::parse(done.outcome.stats, nullptr, false) : nlohmann::json(); const auto counts = CountsOf(stats); if (counts) Database::Get()->SetUgcModelStats(done.id, counts->bricks, counts->triangles, counts->trianglesBefore); // The make with its options and times, for comparing the options on the dashboard if (done.outcome.ok && !done.outcome.options.empty()) { const auto& ms = stats.is_object() ? stats.value("ms", nlohmann::json::object()) : nlohmann::json::object(); const auto part = [&ms](const char* name) { return ms.is_object() ? ms.value(name, 0u) : 0u; }; IUgc::ProcessRun run{ done.id, done.outcome.options, cost.milliseconds, cost.cpuMilliseconds, part("hiddenSurfaces"), part("ambientOcclusion"), part("icon") }; if (counts) { run.bricks = counts->bricks; run.trianglesBefore = counts->trianglesBefore; run.triangles = counts->triangles; } Database::Get()->RecordUgcModelRun(run); } } else { Database::Get()->SetModularBuildProcessed(done.id, state, attempts, error); if (done.outcome.ok) Database::Get()->SetModularBuildProcessStats(done.id, cost); // Which combination's files it shares, for the worlds' manifest answers if (const auto combo = m_ComboOf.find(done.id); combo != m_ComboOf.end()) Database::Get()->SetModularBuildCombination(done.id, combo->second); } m_Recent.erase({ done.kind, done.id }); if (done.outcome.empty) { m_Empty++; LOG_DEBUG("%s %llu has no bricks: nothing to make", KindName(done.kind), static_cast(done.id)); } else if (done.outcome.ok) { m_Made++; m_StoredBytes += done.bytes; LOG_DEBUG("Made %s %llu in %.0f ms%s%s", KindName(done.kind), static_cast(done.id), done.milliseconds, done.outcome.note.empty() ? "" : ": ", done.outcome.note.c_str()); } else { m_Failed++; LOG("Couldn't make %s %llu (attempt %u): %s", KindName(done.kind), static_cast(done.id), attempts, error.c_str()); } m_Log.push_back({ done.kind, done.id, done.outcome.ok || done.outcome.empty, done.milliseconds, done.outcome.empty ? std::string("no bricks: nothing to make") : done.outcome.ok ? done.outcome.note : error, UnixNow() }); while (m_Log.size() > LOG_LENGTH) m_Log.pop_front(); } void UgcProcessor::Collect() { std::deque finished; { std::lock_guard lock(m_Mutex); finished.swap(m_Done); } for (const auto& done : finished) { // done.id is the model, or the combination if (done.outcome.ok) StoreChecksums(done.kind, done.id, done.checksums); if (done.kind == Kind::MODEL) { m_InFlight.erase({ done.kind, done.id }); if (!done.iconOnly) { Record(done); continue; } if (done.outcome.ok && done.iconChangeMs != 0.0) RecordIconTime(done.id, done.iconChangeMs); m_Log.push_back({ Kind::MODEL, done.id, done.outcome.ok, done.milliseconds, done.outcome.ok ? "icon drawn again" : done.outcome.error, UnixNow() }); while (m_Log.size() > LOG_LENGTH) m_Log.pop_front(); continue; } // A combination: every build waiting for it gets its outcome (none when only its icon was drawn again) m_ComboJobs.erase(done.id); if (done.outcome.ok) m_StoredBytes += done.bytes; auto rows = std::move(m_ComboRows[done.id]); m_ComboRows.erase(done.id); for (const auto& [row, attempts] : rows) { m_InFlight.erase({ Kind::MODULAR, row }); Done forRow = done; forRow.id = row; forRow.attempts = attempts; forRow.bytes = 0; Record(forRow); } if (rows.empty()) { m_Log.push_back({ Kind::MODULAR, done.id, done.outcome.ok, done.milliseconds, done.outcome.ok ? "icon drawn again" : done.outcome.error, UnixNow() }); while (m_Log.size() > LOG_LENGTH) m_Log.pop_front(); } } // Something finished: look for more now rather than at the next interval if (!finished.empty()) m_NextPoll = std::min(m_NextPoll, std::chrono::steady_clock::now() + std::chrono::milliseconds(100)); } void UgcProcessor::StoreChecksums(Kind kind, LWOOBJID storageId, const std::vector& checksums) { const auto owner = kind == Kind::MODEL ? IUgc::eFileOwner::MODEL : IUgc::eFileOwner::COMBINATION; for (const auto& checksum : checksums) { // A player model's mesh that changed: the worlds showing it tell their clients (UGC_MODELS_MADE) if (kind == Kind::MODEL && checksum.file == "model.nif") { const auto before = Database::Get()->GetUgcFileChecksum(storageId, checksum.file); if (!before || before->md5 != checksum.md5 || before->size != checksum.size) m_ChangedMeshes.push_back(storageId); } Database::Get()->SetUgcFileChecksum(owner, storageId, checksum.file, checksum.md5, checksum.size); } } void UgcProcessor::Backfill() { if (!m_BackfillBuildsDone) { const auto builds = Database::Get()->GetModularBuildsWithoutCombination(BACKFILL_BUILDS_PER_UPDATE); for (const auto& build : builds) { const auto key = UgcModularKey::Normalize(build.modules); const auto combo = key.empty() ? NO_COMBINATION : UgcModularKey::StorageId(key); m_ComboOf[build.id] = combo; Database::Get()->SetModularBuildCombination(build.id, combo); } if (builds.size() < BACKFILL_BUILDS_PER_UPDATE) m_BackfillBuildsDone = true; } if (!m_StatsBackfillDone) { for (size_t i = 0; i < BACKFILL_ITEMS_PER_UPDATE && !m_StatsBackfill.empty(); i++) { const auto id = m_StatsBackfill.front(); m_StatsBackfill.pop_front(); const auto file = m_Storage.File(Kind::MODEL, id, "stats.json"); if (!file) continue; std::ifstream in(*file, std::ios::binary); const auto counts = CountsOf(nlohmann::json::parse(std::string(std::istreambuf_iterator(in), {}), nullptr, false)); if (counts) Database::Get()->SetUgcModelStats(id, counts->bricks, counts->triangles, counts->trianglesBefore); } if (m_StatsBackfill.empty()) { m_StatsBackfillDone = true; std::ofstream(m_Storage.GetRoot() / STATS_BACKFILL_MARKER) << "1\n"; LOG("Stored the triangle counts of the models made before"); } } if (m_BackfillItemsDone) return; const auto read = [](const std::filesystem::path& path) { std::ifstream in(path, std::ios::binary); return std::string(std::istreambuf_iterator(in), {}); }; for (size_t i = 0; i < BACKFILL_ITEMS_PER_UPDATE && !m_BackfillItems.empty(); i++) { const auto entry = m_BackfillItems.front(); m_BackfillItems.pop_front(); if (m_InFlight.contains({ entry.kind, entry.id }) || (entry.kind == Kind::MODULAR && m_ComboJobs.contains(entry.id))) continue; // Only the icon: it is small (a model's .nif isn't asked for without 3D services and gets its sd0 when made again) const auto checksumFile = m_Storage.File(entry.kind, entry.id, "icon.dds.checksum"); if (!checksumFile) continue; Checksum checksum{ "icon.dds" }; if (!UgcFormats::ReadChecksumXml(read(*checksumFile), checksum.md5, checksum.size)) continue; if (!m_Storage.File(entry.kind, entry.id, "icon.dds.sd0")) { const auto packed = m_Storage.File(entry.kind, entry.id, "icon.dds.gz"); const auto icon = packed ? ZCompression::Gunzip(read(*packed)) : std::nullopt; if (!icon || UgcFormats::Md5Hex(*icon) != checksum.md5) continue; std::string error; if (!m_Storage.Update(entry.kind, entry.id, { { "icon.dds.sd0", Sd0::Compress(*icon) } }, error)) { LOG("Couldn't write the sd0 icon of %llu: %s", static_cast(entry.id), error.c_str()); continue; } } StoreChecksums(entry.kind, entry.id, { checksum }); } if (m_BackfillItems.empty()) { m_BackfillItemsDone = true; std::ofstream(m_Storage.GetRoot() / BACKFILL_MARKER) << "1\n"; LOG("Stored the checksums of the icons made before"); } } void UgcProcessor::Update() { Collect(); // Live update: only the running jobs' outcomes are recorded if (m_Draining) return; Backfill(); const auto now = std::chrono::steady_clock::now(); if (now - m_CpuSampled >= std::chrono::seconds(2)) SampleUsage(); if (now >= m_NextPoll) { m_NextPoll = now + std::chrono::milliseconds(m_Config.pollIntervalMs); Poll(); } 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). 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; }); } LWOOBJID UgcProcessor::StorageId(Kind kind, LWOOBJID id) { if (kind == Kind::MODEL) return id; if (const auto it = m_ComboOf.find(id); it != m_ComboOf.end()) return it->second; const auto info = Database::Get()->GetModularBuildProcessInfo(id); const auto key = info ? UgcModularKey::Normalize(info->details) : std::string(); const LWOOBJID combo = key.empty() ? 0 : UgcModularKey::StorageId(key); if (combo != 0) m_ComboOf[id] = combo; return combo; } UgcProcessor::Availability UgcProcessor::Request(Kind kind, LWOOBJID id) { if (m_InFlight.contains({ kind, id })) return Availability::QUEUED; const auto storageId = StorageId(kind, id); if (kind == Kind::MODULAR && m_ComboJobs.contains(storageId)) return Availability::QUEUED; if (storageId != 0 && m_Storage.File(kind, storageId, "icon.png")) { m_Storage.Touch(kind, storageId); return Availability::READY; } if (const auto it = m_Recent.find({ kind, id }); it != m_Recent.end()) return it->second.second; const auto info = kind == Kind::MODEL ? Database::Get()->GetUgcProcessInfo(id) : Database::Get()->GetModularBuildProcessInfo(id); Availability answer = Availability::UNKNOWN; // Failed, or nothing to make (no bricks): 404, as for HKX, so the client doesn't wait if (info && info->state != IUgc::eProcessState::FAILED && info->state != IUgc::eProcessState::EMPTY) { // Made before but the files are gone (deleted to save space): make them again, first if (info->state == IUgc::eProcessState::DONE) { if (kind == Kind::MODEL) Database::Get()->ResetUgcModelProcessing(id, false, ""); else Database::Get()->ResetModularBuildProcessing(id, false); } // Someone wants it now: no need to wait out the quiet period after its save if (kind == Kind::MODEL) Database::Get()->ExpediteUgcModel(id); m_NextPoll = std::chrono::steady_clock::now(); answer = Availability::QUEUED; } m_Recent[{ kind, id }] = { std::chrono::steady_clock::now(), answer }; return answer; } nlohmann::json UgcProcessor::Status() const { nlohmann::json status; { std::lock_guard lock(m_Mutex); status["queued"] = m_Jobs.size(); status["purge"] = PurgeStatus(); status["active"] = m_Active; } Limits limits; { std::lock_guard lock(m_Mutex); limits = m_Limits; status["jobMemoryBytes"] = m_MemoryInUse; status["memoryWaits"] = m_Waiting; } status["workers"] = m_Threads.size(); status["reused"] = m_Reused; const auto throttle = UgcThrottle::GetStats(); status["usage"] = { { "cpuPercent", std::lround(m_CpuPercent) }, // of one core { "cores", std::thread::hardware_concurrency() }, { "residentBytes", ResidentBytes() }, { "throttled", Throttled() }, { "throttledMs", throttle.sleptMs }, { "paused", m_Paused }, }; status["limits"] = { { "maxCpus", limits.maxCpus }, { "maxMemoryBytes", limits.maxMemoryBytes }, { "nice", limits.nice }, { "pauseHours", limits.pauseFromHour >= 0 ? std::to_string(limits.pauseFromHour) + "-" + std::to_string(limits.pauseToHour) : "" }, }; status["made"] = m_Made; status["failed"] = m_Failed; status["empty"] = m_Empty; status["evicted"] = m_Evicted; status["storedBytes"] = m_StoredBytes; status["maxStorageBytes"] = m_Config.maxStorageBytes; auto& recent = status["recent"] = nlohmann::json::array(); for (auto it = m_Log.rbegin(); it != m_Log.rend(); ++it) { recent.push_back({ { "kind", KindName(it->kind) }, { "id", std::to_string(it->id) }, { "ok", it->ok }, { "ms", static_cast(it->milliseconds) }, { "message", it->message }, { "time", it->time } }); } return status; } UgcIconParams::Values UgcProcessor::IconValues(const std::string& kind, const std::string& itemTarget) { UgcIconParams::Values values; // Player models have no shared preset: every one is a different size and shape, so each is fitted to the icon from // the defaults, and only its own settings change it. Cars and rockets of a build type share a preset. if (kind != UgcIconParams::ModelKind()) { if (const auto preset = Database::Get()->GetUgcIconSettings(UgcIconParams::KindTarget(kind))) values = UgcIconParams::Parse(*preset); } if (const auto own = Database::Get()->GetUgcIconSettings(itemTarget)) { for (const auto& [key, value] : UgcIconParams::Parse(*own)) values[key] = value; } return values; } bool UgcProcessor::QueuePreview(Kind kind, LWOOBJID id, const std::string& modules, const UgcIconParams::Values& values, DeferredReply reply, std::string& error) { Job job{ kind, id, 0 }; if (kind == Kind::MODULAR) { if (!UgcCdClient::GatherModular(modules, job.modular, error)) return false; job.modular.key = UgcModularKey::Normalize(modules); job.modular.iconValues = values; } else { if (!m_Storage.File(Kind::MODEL, id, "model.nif.gz") && !m_Storage.File(Kind::MODEL, id, "model.nif")) { error = "the model has no stored .nif yet"; return false; } job.iconValues = values; } job.preview = std::move(reply); { std::lock_guard lock(m_Mutex); job.memory = UgcJobs::EstimateMemory(64, m_Settings); m_Jobs.push_front(std::move(job)); } m_Wake.notify_all(); return true; } std::shared_ptr UgcProcessor::CachedAssembly(const std::string& key) { std::lock_guard lock(m_AssemblyMutex); const auto it = std::find_if(m_Assemblies.begin(), m_Assemblies.end(), [&key](const auto& entry) { return entry.first == key; }); if (it == m_Assemblies.end()) return nullptr; m_Assemblies.splice(m_Assemblies.begin(), m_Assemblies, it); return m_Assemblies.front().second; } void UgcProcessor::CacheAssembly(const std::string& key, std::shared_ptr nif) { constexpr size_t MAX_ENTRIES = 32; constexpr size_t MAX_BYTES = 64ull * 1024 * 1024; std::lock_guard lock(m_AssemblyMutex); std::erase_if(m_Assemblies, [&key](const auto& entry) { return entry.first == key; }); m_Assemblies.emplace_front(key, std::move(nif)); size_t bytes = 0, kept = 0; for (auto it = m_Assemblies.begin(); it != m_Assemblies.end(); ++it, kept++) { bytes += it->second->size(); if (kept >= MAX_ENTRIES || (kept > 0 && bytes > MAX_BYTES)) { m_Assemblies.erase(it, m_Assemblies.end()); break; } } } bool UgcProcessor::QueueAssembly(const std::string& modules, DeferredReply reply, std::string& error) { const auto key = UgcModularKey::Normalize(modules); if (key.empty()) { error = "no modules"; return false; } if (const auto cached = CachedAssembly(key)) { HTTPReply out; out.status = eHTTPStatusCode::OK; out.contentType = eContentType::APPLICATION_OCTET_STREAM; out.message = *cached; out.headers.push_back("Cache-Control: no-store"); reply.Send(std::move(out)); return true; } Job job{ Kind::MODULAR, 0, 0 }; if (!UgcCdClient::GatherModular(modules, job.modular, error)) return false; job.modular.key = key; job.assembly = true; job.preview = std::move(reply); { std::lock_guard lock(m_Mutex); job.memory = UgcJobs::EstimateMemory(64, m_Settings); m_Jobs.push_front(std::move(job)); } m_Wake.notify_all(); return true; } size_t UgcProcessor::RegenerateIcons(const std::string& kind) { std::vector jobs; UgcJobs::Settings settings; { std::lock_guard lock(m_Mutex); settings = m_Settings; } const bool models = kind == UgcIconParams::ModelKind(); for (const auto& entry : m_Storage.List()) { if (models) { if (entry.kind != Kind::MODEL || m_InFlight.contains({ Kind::MODEL, entry.id }) || (!m_Storage.File(Kind::MODEL, entry.id, "model.nif.gz") && !m_Storage.File(Kind::MODEL, entry.id, "model.nif"))) continue; Job job{ Kind::MODEL, entry.id, 0 }; job.iconOnly = true; job.iconValues = IconValues(kind, UgcIconParams::ModelTarget(entry.id)); job.memory = UgcJobs::EstimateMemory(64, settings); m_InFlight.insert({ Kind::MODEL, entry.id }); jobs.push_back(std::move(job)); continue; } if (entry.kind != Kind::MODULAR || m_ComboJobs.contains(entry.id)) continue; const auto path = m_Storage.File(Kind::MODULAR, entry.id, "combo.json"); const auto text = path ? UgcBricks::ReadFile(*path) : std::nullopt; const auto combo = text ? nlohmann::json::parse(*text, nullptr, false) : nlohmann::json(); if (!combo.is_object() || UgcIconParams::BuildKind(combo.value("buildType", -1)) != kind) continue; const auto key = combo.value("key", std::string()); Job job{ Kind::MODULAR, entry.id, 0 }; std::string error; // The key's LOTs written as an ldf_config ("4713-4714" -> "4713+4714") std::string modules = key; std::replace(modules.begin(), modules.end(), '-', '+'); if (key.empty() || !UgcCdClient::GatherModular(modules, job.modular, error)) continue; job.modular.key = key; job.modular.iconValues = IconValues(kind, UgcIconParams::CombinationTarget(key)); job.memory = UgcJobs::EstimateMemory(64, settings); m_ComboJobs.insert(entry.id); jobs.push_back(std::move(job)); } const auto count = jobs.size(); { std::lock_guard lock(m_Mutex); for (auto& job : jobs) m_Jobs.push_back(std::move(job)); } m_Wake.notify_all(); return count; } UgcProcessor::DeleteResult UgcProcessor::Delete(const DeleteRequest& request) { DeleteResult result; 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) m_PurgeRowsOf[storageId].push_back(id); } 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; } 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)); }