mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-02 10:53:44 +00:00
- IUgc::GetUgcProcessTotals sums what the UGC server made, for models and for car and rocket builds: count, time spent (and how many are timed), CPU time, average and slowest, average and most memory (estimate), bricks, triangles and triangles saved (models whose count before hidden face removal is known). /api/ugc returns them as totals and the UGC page shows a card for each kind. Parity tested. - Player models no longer have a shared icon preset: every one is a different size and shape, so its icon is fitted to it from the settings. The UGC server ignores a kind:model preset, the dashboard refuses to save one, and the icon editor hides the type buttons for models. One model's own icon values still work. Car and rocket build types keep theirs. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1002 lines
41 KiB
C++
1002 lines
41 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;
|
|
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();
|
|
}
|
|
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 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 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;
|
|
UgcThrottle::Checkpoint();
|
|
|
|
{
|
|
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));
|
|
}
|