Files
DarkflameServer/dUgcServer/Processing/UgcProcessor.cpp
Aaron Kimbrell 7c32779e89 fix(ugc): transparent bricks drawn see-through in mixed models; dUgcServer in folders by role
- The client puts a shape in its sorted, blended pass only when its
  NiMaterialProperty alpha is under 0.99999 (ShaderCommon::GetAlphaFlags
  0x0109f5a0; the NiAlphaProperty blend flag isn't read); at 1.0 it's drawn
  solid with blending off. Transparent (and transparent glitter) shapes now
  get a material with alpha 0.9999, as the S01_Alpha shapes of the game's own
  brick models (res/BrickModels/ndmade) do; opaque shapes keep 1.0. Models
  with a transparent brick change; the others are byte for byte the same.
- dUgcServer's files move into Bricks/, Model/, Render/, Formats/ and
  Processing/ (the CMakeLists says what each holds); includes are unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-28 22:31:26 -05:00

1020 lines
42 KiB
C++

#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 <algorithm>
#include <ctime>
#include <filesystem>
#include <fstream>
#include <thread>
#if defined(__linux__)
#include <malloc.h>
#include <sys/resource.h>
#include <sys/syscall.h>
#include <unistd.h>
#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<ModelCounts> 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<double>(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<id_t>(syscall(SYS_gettid)), std::clamp(nice, 0, 19));
#else
(void)nice;
#endif
}
int64_t UnixNow() {
return std::chrono::duration_cast<std::chrono::seconds>(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<uint64_t>(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<double>(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<size_t>(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<unsigned long long>(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<Job> 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<const std::string>(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);
outcome.ok = nif && UgcJobs::IconFromNif(*nif, options, outcome.files, outcome.error, settings.shaders.TagLooks());
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::ThreadCpuSeconds();
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);
done.outcome.ok = nif && UgcJobs::IconFromNif(*nif, options, done.outcome.files, done.outcome.error, settings.shaders.TagLooks());
if (!nif) done.outcome.error = "no stored .nif";
} else {
done.outcome = job.kind == Kind::MODEL
? UgcJobs::ProcessModel(job.blob, m_Library, settings, static_cast<uint64_t>(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<double, std::milli>(std::chrono::steady_clock::now() - start).count();
done.cpuMilliseconds = std::max(0.0, UgcThrottle::ThreadCpuSeconds() - 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<size_t>(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<uint32_t>(wanted - queued - active + m_InFlight.size());
const auto buildLimit = std::max<uint32_t>(limit, static_cast<uint32_t>(m_Threads.size() + m_InFlight.size()));
std::vector<Job> jobs;
std::vector<IUgc::PendingModel> 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;
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::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<uint32_t>(done.milliseconds), static_cast<uint32_t>(done.cpuMilliseconds), static_cast<uint32_t>(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();
if (const auto counts = CountsOf(stats)) Database::Get()->SetUgcModelStats(done.id, counts->bricks, counts->triangles, counts->trianglesBefore);
} 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<unsigned long long>(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<unsigned long long>(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<unsigned long long>(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<Done> 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;
}
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<Checksum>& 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<char>(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<char>(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<unsigned long long>(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<int64_t>(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<const std::string> 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<const std::string> 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<Job> 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<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));
}