mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-02 19:03:43 +00:00
- 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>
1020 lines
42 KiB
C++
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));
|
|
}
|