diff --git a/dDashboardServer/DashboardServer.cpp b/dDashboardServer/DashboardServer.cpp index 19048bfcb..1049b138e 100644 --- a/dDashboardServer/DashboardServer.cpp +++ b/dDashboardServer/DashboardServer.cpp @@ -654,7 +654,9 @@ int main(int argc, char** argv) { } - // Cleanup + // Cleanup: the conversion threads first (they answer deferred requests), then the web server's connections + Scenery::Shutdown(); + Game::web.Shutdown(); Inspector::Shutdown(); EmailService::Shutdown(); ModeratorHelper::Shutdown(); diff --git a/dDashboardServer/routes/CMakeLists.txt b/dDashboardServer/routes/CMakeLists.txt index 8548fac80..fd3ffe32a 100644 --- a/dDashboardServer/routes/CMakeLists.txt +++ b/dDashboardServer/routes/CMakeLists.txt @@ -42,6 +42,7 @@ set(DASHBOARDROUTES_SOURCES "WorldView.cpp" "NifFile.cpp" "Scenery.cpp" + "WorkerPool.cpp" "AccountModeration.cpp" "ModerationTools.cpp" "ModeratorHelper.cpp" diff --git a/dDashboardServer/routes/ClientAssets.cpp b/dDashboardServer/routes/ClientAssets.cpp index 4d44a10fb..824d1eeb0 100644 --- a/dDashboardServer/routes/ClientAssets.cpp +++ b/dDashboardServer/routes/ClientAssets.cpp @@ -47,8 +47,8 @@ namespace { * Where a normalized (lowercase) res/ path is on disk. Unpacked clients keep their original mixed case, so on * case-sensitive file systems each part is matched ignoring case. */ - std::filesystem::path ResolveRes(const std::string& normalized) { - auto current = ClientRes(); + std::filesystem::path ResolveResIn(const std::filesystem::path& res, const std::string& normalized) { + auto current = res; std::error_code ec; for (const auto& part : GeneralUtils::SplitString(normalized, '/')) { if (part.empty()) continue; @@ -72,6 +72,10 @@ namespace { return current; } + std::filesystem::path ResolveRes(const std::string& normalized) { + return ResolveResIn(ClientRes(), normalized); + } + std::optional ReadFile(const std::filesystem::path& path) { std::ifstream file(path, std::ios::binary | std::ios::ate); if (!file) return std::nullopt; @@ -338,6 +342,19 @@ namespace ClientAssets { return path; } + std::filesystem::path ResFolder() { + return ClientRes(); + } + + std::optional ResolveResFile(const std::string& relativePath, const std::filesystem::path& res) { + const auto normalized = NormalizeAssetPath(relativePath); + if (res.empty() || !IsSafeAssetPath(normalized)) return std::nullopt; + auto path = ResolveResIn(res, normalized); + std::error_code ec; + if (!std::filesystem::is_regular_file(path, ec)) return std::nullopt; + return path; + } + std::optional FindResFile(const std::string& folder, const std::string& fileName) { const auto normalized = NormalizeAssetPath(folder); if (ClientRes().empty() || !IsSafeAssetPath(normalized)) return std::nullopt; diff --git a/dDashboardServer/routes/ClientAssets.h b/dDashboardServer/routes/ClientAssets.h index be800d9c2..79b264b80 100644 --- a/dDashboardServer/routes/ClientAssets.h +++ b/dDashboardServer/routes/ClientAssets.h @@ -37,6 +37,12 @@ namespace ClientAssets { // Where a file inside res/ is on disk (matched ignoring case), if it exists std::optional ResolveResFile(const std::string& relativePath); + // The client's res/ folder (empty when client_location isn't set) + std::filesystem::path ResFolder(); + + // ResolveResFile inside a res/ folder from ResFolder: reads no settings, so worker threads may call it + std::optional ResolveResFile(const std::string& relativePath, const std::filesystem::path& res); + // The res/-relative path of a file called `fileName` (ignoring case) anywhere under res/ std::optional FindResFile(const std::string& folder, const std::string& fileName); diff --git a/dDashboardServer/routes/Scenery.cpp b/dDashboardServer/routes/Scenery.cpp index a48474b8f..5913a20ca 100644 --- a/dDashboardServer/routes/Scenery.cpp +++ b/dDashboardServer/routes/Scenery.cpp @@ -5,9 +5,13 @@ #include #include #include +#include #include #include +#include #include +#include +#include #include #include @@ -16,6 +20,8 @@ #include "ReportRoutes.h" #include "RouteUtils.h" #include "TtlCache.h" +#include "Web.h" +#include "WorkerPool.h" #include "WorldScene.h" #include "ZonePaths.h" @@ -194,6 +200,9 @@ namespace { std::map index; // res path -> its index in assets std::string json; std::optional flairs; // the flairs' manifest, built when first asked for (their models join assets) + std::unordered_set flairModels; // res paths of the flairs' models, converted ahead of others + bool warmedScenery{}; // WarmUp queued the scenery's models + bool warmedFlairs{}; // and the flairs' size_t IndexOf(const std::string& path) { const auto [it, added] = index.try_emplace(path, assets.size()); @@ -277,6 +286,7 @@ namespace { const auto model = models.find(flair.id); if (model == models.end() || model->second.empty() || !std::isfinite(flair.position.x)) continue; assetOf.push_back(scenery.IndexOf(model->second)); + scenery.flairModels.insert(model->second); for (const auto value : { flair.position.x, flair.position.y, flair.position.z }) positions.push_back(Round(value, 100.0)); // Radians about x, y and z, applied in that order const double c1 = std::cos(flair.rotation.x / 2), c2 = std::cos(flair.rotation.y / 2), c3 = std::cos(flair.rotation.z / 2); @@ -316,74 +326,170 @@ namespace { return data; } - // Keeps the disk cache under DISK_CACHE_BYTES by removing the least recently written files - void StoreOnDisk(const std::filesystem::path& target, const std::string& data) { - static std::optional total; - std::error_code ec; - std::filesystem::create_directories(CACHE_DIR, ec); - if (!total) { - total = 0; - for (const auto& entry : std::filesystem::directory_iterator(CACHE_DIR, ec)) *total += entry.is_regular_file(ec) ? entry.file_size(ec) : 0; + /** + * The converted models on disk (CACHE_DIR), kept under DISK_CACHE_BYTES by removing the least recently written + * files. Any thread. + */ + class DiskCache { + public: + void Store(const std::filesystem::path& target, const std::string& data) { + std::lock_guard lock(m_Mutex); + std::error_code ec; + std::filesystem::create_directories(CACHE_DIR, ec); + CountLocked(); + if (m_Total + data.size() > DISK_CACHE_BYTES) { + std::vector> files; + for (const auto& entry : std::filesystem::directory_iterator(CACHE_DIR, ec)) { + if (entry.path().extension() == ".bin") files.emplace_back(entry.last_write_time(ec), entry.path()); + } + std::sort(files.begin(), files.end()); + for (const auto& [time, path] : files) { + if (m_Total + data.size() <= DISK_CACHE_BYTES * 3 / 4) break; + const auto size = std::filesystem::file_size(path, ec); + if (std::filesystem::remove(path, ec)) m_Total -= std::min(m_Total, size); + } + } + // Written next to the target and renamed, so a reader never sees half a file + std::ostringstream temporary; + temporary << target.string() << "." << std::this_thread::get_id() << ".tmp"; + { + std::ofstream file(temporary.str(), std::ios::binary | std::ios::trunc); + if (!file.write(data.data(), static_cast(data.size()))) return; + } + const auto existed = std::filesystem::exists(target, ec); + std::filesystem::rename(temporary.str(), target, ec); + if (!ec && !existed) m_Total += data.size(); } - if (*total + data.size() > DISK_CACHE_BYTES) { - std::vector> files; - for (const auto& entry : std::filesystem::directory_iterator(CACHE_DIR, ec)) files.emplace_back(entry.last_write_time(ec), entry.path()); - std::sort(files.begin(), files.end()); - for (const auto& [time, path] : files) { - if (*total + data.size() <= DISK_CACHE_BYTES * 3 / 4) break; - const auto size = std::filesystem::file_size(path, ec); - if (std::filesystem::remove(path, ec)) *total -= std::min(*total, size); + + // Bytes in the cache + uintmax_t Total() { + std::lock_guard lock(m_Mutex); + CountLocked(); + return m_Total; + } + + private: + void CountLocked() { + if (m_Counted) return; + m_Counted = true; + std::error_code ec; + for (const auto& entry : std::filesystem::directory_iterator(CACHE_DIR, ec)) { + if (entry.path().extension() == ".tmp") std::filesystem::remove(entry.path(), ec); // left by a crash + else m_Total += entry.is_regular_file(ec) ? entry.file_size(ec) : 0; } } - const auto temporary = target.string() + ".tmp"; - { - std::ofstream file(temporary, std::ios::binary | std::ios::trunc); - if (!file.write(data.data(), static_cast(data.size()))) return; + + std::mutex m_Mutex; + uintmax_t m_Total{}; + bool m_Counted{}; + }; + + DiskCache g_Disk; + + using Bytes = std::shared_ptr; + + // A TtlCache any thread may use + class SharedCache { + public: + SharedCache(std::chrono::seconds ttl, size_t maxBytes) : m_Cache(ttl, maxBytes) {} + Bytes Get(const std::string& key) { + std::lock_guard lock(m_Mutex); + const auto cached = m_Cache.Get(key); + return cached ? *cached : nullptr; } - std::filesystem::rename(temporary, target, ec); - if (!ec) *total += data.size(); + void Put(const std::string& key, Bytes value) { + if (!value) return; + std::lock_guard lock(m_Mutex); + const auto weight = value->size() + 256; + m_Cache.Put(key, std::move(value), weight); + } + private: + std::mutex m_Mutex; + TtlCache m_Cache; + }; + + SharedCache g_Models(std::chrono::hours(1), MESH_CACHE_BYTES); // "path|lod" -> NifFile::Encode's output + SharedCache g_Embedded(std::chrono::hours(1), MESH_CACHE_BYTES); // "path#block" -> DDS of a texture stored in a .nif + + // Conversions under way ("path|lod"), so a model asked for twice at once is converted once and both get it + std::mutex g_ConvertingMutex; + std::map> g_Converting; + + std::string ModelKey(const std::string& path, uint32_t lod) { return path + "|" + std::to_string(lod); } + + // Where model `path` at `lod` is kept on disk, named after the source file's size and time so a changed client + // file is converted again + std::filesystem::path DiskPath(const std::string& key, const std::filesystem::path& file) { + std::error_code ec; + const auto size = std::filesystem::file_size(file, ec); + const auto time = std::filesystem::last_write_time(file, ec).time_since_epoch().count(); + const auto diskKey = key + "|" + std::to_string(size) + "|" + std::to_string(time) + "|" + std::to_string(FORMAT_VERSION); + return CACHE_DIR / (std::to_string(Fnv1a(diskKey)) + ".bin"); + } + + Bytes Convert(const std::string& path, uint32_t lod, const std::filesystem::path& file, const std::filesystem::path& target) { + if (auto encoded = ReadWhole(target)) return std::make_shared(std::move(*encoded)); + const auto data = ReadWhole(file); + if (!data) return nullptr; + std::string error; + const auto model = NifFile::Parse(*data, lod, error); + if (!model) { + LOG_DEBUG("Could not read %s: %s", path.c_str(), error.c_str()); + return nullptr; + } + const auto folder = FolderOf(path); + std::vector textures; // per mesh: a res path, "#" for one stored in the .nif, or empty + for (const auto& mesh : model->meshes) { + if (mesh.material.embeddedTexture >= 0) textures.push_back("#" + std::to_string(mesh.material.embeddedTexture)); + else textures.push_back(mesh.material.texture.empty() ? std::string{} : FindTexture(folder, mesh.material.texture)); + } + auto encoded = std::make_shared(NifFile::Encode(*model, textures)); + g_Disk.Store(target, *encoded); + return encoded; } /** - * Model `path` at `lod` in NifFile::Encode's format. Converting a big .nif takes a moment, so results are kept in - * memory (MESH_CACHE_BYTES) and on disk (DISK_CACHE_BYTES), keyed by the file's size and time so a changed client - * file is converted again. + * Model `path` (on disk at `file`) at `lod` in NifFile::Encode's format. Converting a big .nif takes a moment, so + * results are kept in memory (MESH_CACHE_BYTES; unless `keep` is false, for conversions ahead of time) and on disk + * (DISK_CACHE_BYTES). Any thread; the same model asked for again while it converts waits for that conversion. */ - std::shared_ptr Encoded(const std::string& path, uint32_t lod) { - static TtlCache> memory(std::chrono::hours(1), MESH_CACHE_BYTES); - const auto key = path + "|" + std::to_string(lod); - if (auto cached = memory.Get(key)) return *cached; + Bytes Encoded(const std::string& path, uint32_t lod, const std::filesystem::path& file, bool keep = true) { + const auto key = ModelKey(path, lod); + if (auto cached = g_Models.Get(key)) return cached; - const auto file = ClientAssets::ResolveResFile(path); - if (!file) return nullptr; - std::error_code ec; - const auto size = std::filesystem::file_size(*file, ec); - const auto time = std::filesystem::last_write_time(*file, ec).time_since_epoch().count(); - const auto diskKey = key + "|" + std::to_string(size) + "|" + std::to_string(time) + "|" + std::to_string(FORMAT_VERSION); - const auto target = CACHE_DIR / (std::to_string(Fnv1a(diskKey)) + ".bin"); - - auto encoded = ReadWhole(target); - if (!encoded) { - const auto data = ReadWhole(*file); - if (!data) return nullptr; - std::string error; - const auto model = NifFile::Parse(*data, lod, error); - if (!model) { - LOG_DEBUG("Could not read %s: %s", path.c_str(), error.c_str()); - return nullptr; + std::promise promise; + std::shared_future converting; + bool mine = false; + { + std::lock_guard lock(g_ConvertingMutex); + const auto it = g_Converting.find(key); + if (it != g_Converting.end()) { + converting = it->second; + } else { + converting = promise.get_future().share(); + g_Converting.emplace(key, converting); + mine = true; } - const auto folder = FolderOf(path); - std::vector textures; // per mesh: a res path, "#" for one stored in the .nif, or empty - for (const auto& mesh : model->meshes) { - if (mesh.material.embeddedTexture >= 0) textures.push_back("#" + std::to_string(mesh.material.embeddedTexture)); - else textures.push_back(mesh.material.texture.empty() ? std::string{} : FindTexture(folder, mesh.material.texture)); - } - encoded = NifFile::Encode(*model, textures); - StoreOnDisk(target, *encoded); } - auto shared = std::make_shared(std::move(*encoded)); - memory.Put(key, shared, shared->size() + 256); - return shared; + if (!mine) { + auto result = converting.get(); + if (keep) g_Models.Put(key, result); + return result; + } + + Bytes result; + try { + result = Convert(path, lod, file, DiskPath(key, file)); + } catch (const std::exception& ex) { + LOG("Could not convert %s: %s", path.c_str(), ex.what()); + } + if (keep) g_Models.Put(key, result); + { + std::lock_guard lock(g_ConvertingMutex); + g_Converting.erase(key); + } + promise.set_value(result); + return result; } // The "textures" list of an encoded model's header @@ -409,12 +515,145 @@ namespace { uint32_t LodOf(const HTTPContext& context) { return std::min(GeneralUtils::TryParse(QueryValue(context.queryString, "lod")).value_or(0), MAX_LOD); } + + // ---- Converting on worker threads, so a big model never holds up the web server's one thread ---- + + WorkerPool g_Pool; + + constexpr uintmax_t SMALL_MODEL_BYTES = 256 * 1024; // .nif files this small convert in the pool's fast lane + constexpr uintmax_t LARGE_MODEL_BYTES = 4 * 1024 * 1024; // and this big wait behind everything smaller + constexpr auto WARM_IDLE = std::chrono::seconds(90); // converting a zone ahead stops once nobody has asked for it this long + constexpr uintmax_t WARM_DISK_BYTES = DISK_CACHE_BYTES * 3 / 4; // and when the disk cache is this full (it never evicts for it) + constexpr size_t WARM_MAX_MODELS = 4000; + + // When someone last asked for something of a zone, and the LOD they last asked a model at + struct Activity { + std::chrono::steady_clock::time_point last; + uint32_t lod{ 1 }; // the world view's default detail + }; + std::mutex g_ActivityMutex; + std::map g_Activity; + + void Touch(uint32_t zoneId, std::optional lod = std::nullopt) { + std::lock_guard lock(g_ActivityMutex); + auto& activity = g_Activity[zoneId]; + activity.last = std::chrono::steady_clock::now(); + if (lod) activity.lod = *lod; + } + + // The zone's activity while someone views it, nullopt once nobody has for WARM_IDLE + std::optional Viewed(uint32_t zoneId) { + std::lock_guard lock(g_ActivityMutex); + const auto it = g_Activity.find(zoneId); + if (it == g_Activity.end() || std::chrono::steady_clock::now() - it->second.last > WARM_IDLE) return std::nullopt; + return it->second; + } + + // Flairs (small, and drawn around the camera) and small models first; big ones behind the rest + WorkerPool::ePriority PriorityOf(const ZoneScenery& zone, const std::string& path, const std::filesystem::path& file) { + if (zone.flairModels.contains(path)) return WorkerPool::ePriority::URGENT; + std::error_code ec; + const auto size = std::filesystem::file_size(file, ec); + if (ec || size <= SMALL_MODEL_BYTES) return WorkerPool::ePriority::URGENT; + return size >= LARGE_MODEL_BYTES ? WorkerPool::ePriority::LARGE : WorkerPool::ePriority::NORMAL; + } + + uint64_t WarmGroup(uint32_t zoneId) { return static_cast(zoneId) + 1; } + + /** + * Convert models of a zone ahead of time (onto the disk cache), at the LOD its viewer last asked for, while + * someone views it. The lowest priority: only when nothing else waits. Smallest first; `front` puts these before + * the zone's other queued ones (the flairs). Web thread. + */ + void WarmUp(uint32_t zoneId, const std::vector& paths, bool front) { + if (!g_Pool.Running() || paths.empty()) return; + Files(); // built here: workers only read it + struct Item { + std::string path; + std::filesystem::path file; + uintmax_t size{}; + }; + std::vector items; + std::set seen; + std::error_code ec; + for (const auto& path : paths) { + if (!seen.insert(path).second) continue; + const auto file = ClientAssets::ResolveResFile(path); + if (file) items.push_back({ path, *file, std::filesystem::file_size(*file, ec) }); + } + std::stable_sort(items.begin(), items.end(), [](const Item& a, const Item& b) { return a.size < b.size; }); + if (items.size() > WARM_MAX_MODELS) items.resize(WARM_MAX_MODELS); + const auto group = WarmGroup(zoneId); + // Queued at the front in reverse, so they still run smallest first + if (front) std::reverse(items.begin(), items.end()); + for (auto& item : items) { + g_Pool.Submit(WorkerPool::ePriority::BACKGROUND, [zoneId, group, path = std::move(item.path), file = std::move(item.file)] { + const auto activity = Viewed(zoneId); + if (!activity || g_Disk.Total() >= WARM_DISK_BYTES) { + g_Pool.Cancel(group); + return; + } + const auto key = ModelKey(path, activity->lod); + std::error_code ec; + if (g_Models.Get(key) || std::filesystem::exists(DiskPath(key, file), ec)) return; + Encoded(path, activity->lod, file, false); + }, group, front); + } + } + + // Warm the zone's models when its manifest is asked for (again after nobody viewed it for a while) + void WarmScenery(uint32_t zoneId, ZoneScenery& zone, bool flairs) { + if (!Viewed(zoneId)) zone.warmedScenery = zone.warmedFlairs = false; + Touch(zoneId); + auto& warmed = flairs ? zone.warmedFlairs : zone.warmedScenery; + if (warmed) return; + warmed = true; + if (flairs) WarmUp(zoneId, { zone.flairModels.begin(), zone.flairModels.end() }, true); + else WarmUp(zoneId, zone.assets, false); + } + + // The texture `name` of model `path` (textures[slot] of its encoded form) as a DDS file. Any thread. + std::optional TextureBytes(const std::string& path, const std::filesystem::path& modelFile, const std::vector& textures, + const std::string& name, const std::filesystem::path& res) { + if (!name.starts_with('#')) { + const auto file = ClientAssets::ResolveResFile(name, res); + return file ? ReadWhole(*file) : std::nullopt; + } + // Stored inside the model: reading a big .nif again for each of its textures would be slow, so they're kept + if (const auto cached = g_Embedded.Get(path + name)) return *cached; + const auto data = ReadWhole(modelFile); + const auto block = GeneralUtils::TryParse(name.substr(1)); + if (!data || !block) return std::nullopt; + // Every texture of the file at once, since the browser asks for them together + for (const auto& other : textures) { + if (other == name) continue; + const auto otherBlock = other.starts_with('#') ? GeneralUtils::TryParse(other.substr(1)) : std::nullopt; + auto file = otherBlock ? NifFile::EmbeddedTexture(*data, *otherBlock) : std::nullopt; + if (file) g_Embedded.Put(path + other, std::make_shared(std::move(*file))); + } + auto dds = NifFile::EmbeddedTexture(*data, *block); + if (dds) g_Embedded.Put(path + name, std::make_shared(*dds)); + return dds; + } + + // The model path of `asset` in the zone's manifests, or a 404 reply + const std::string* AssetPath(HTTPReply& reply, uint32_t zoneId, uint32_t asset) { + auto& zone = Zone(zoneId); + // The flairs' models join the list when their manifest is first built (a browser may still have it cached) + if (zone && asset >= zone->assets.size()) Flairs(zoneId, *zone); + if (!zone || asset >= zone->assets.size()) { + JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such model in this zone"); + return nullptr; + } + return &zone->assets[asset]; + } } namespace Scenery { std::optional ZoneJson(uint32_t zoneId) { - const auto& zone = Zone(zoneId); + auto& zone = Zone(zoneId); if (!zone) return std::nullopt; + WarmScenery(zoneId, *zone, false); return zone->json; } @@ -426,61 +665,85 @@ namespace Scenery { std::optional FlairsJson(uint32_t zoneId) { auto& zone = Zone(zoneId); if (!zone) return std::nullopt; - return Flairs(zoneId, *zone); + const auto& flairs = Flairs(zoneId, *zone); + WarmScenery(zoneId, *zone, true); + return flairs; } - void ReplyMesh(HTTPReply& reply, uint32_t zoneId, uint32_t asset, uint32_t lod) { - auto& zone = Zone(zoneId); - // The flairs' models join the list when their manifest is first built (a browser may still have it cached) - if (zone && asset >= zone->assets.size()) Flairs(zoneId, *zone); - if (!zone || asset >= zone->assets.size()) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such model in this zone"); - const auto encoded = Encoded(zone->assets[asset], std::min(lod, MAX_LOD)); - if (!encoded) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "Could not read this model"); - Binary(reply, *encoded); + void ReplyMesh(HTTPReply& reply, const HTTPContext& context, uint32_t zoneId, uint32_t asset, uint32_t lod) { + const auto* found = AssetPath(reply, zoneId, asset); + if (!found) return; + const auto path = *found; + lod = std::min(lod, MAX_LOD); + Touch(zoneId, lod); + if (const auto cached = g_Models.Get(ModelKey(path, lod))) return Binary(reply, *cached); + const auto file = ClientAssets::ResolveResFile(path); + if (!file) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "Could not read this model"); + Files(); // built here: workers only read it + const auto deferred = Web::Defer(reply, context); + g_Pool.Submit(PriorityOf(*Zone(zoneId), path, *file), [deferred, path, lod, file = *file] { + if (deferred.Cancelled()) return; + HTTPReply out; + const auto encoded = Encoded(path, lod, file); + if (encoded) Binary(out, *encoded); + else JsonError(out, eHTTPStatusCode::NOT_FOUND, "Could not read this model"); + deferred.Send(std::move(out)); + }); } - void ReplyTexture(HTTPReply& reply, uint32_t zoneId, uint32_t asset, uint32_t slot, uint32_t lod) { - auto& zone = Zone(zoneId); - if (zone && asset >= zone->assets.size()) Flairs(zoneId, *zone); - if (!zone || asset >= zone->assets.size()) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such model in this zone"); - const auto encoded = Encoded(zone->assets[asset], std::min(lod, MAX_LOD)); - const auto textures = encoded ? TexturesOf(*encoded) : std::vector{}; - if (slot >= textures.size() || textures[slot].empty()) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such texture"); - const auto& texture = textures[slot]; - std::optional dds; - if (texture.starts_with('#')) { - // Stored inside the model: reading a big .nif again for each of its textures would be slow, so they're kept - static TtlCache> embedded(std::chrono::hours(1), MESH_CACHE_BYTES); - const auto key = zone->assets[asset] + texture; - if (const auto cached = embedded.Get(key)) return Binary(reply, **cached); - const auto data = ClientAssets::ReadResFile(zone->assets[asset]); - const auto block = GeneralUtils::TryParse(texture.substr(1)); - if (data && block) { - // Every texture of the file at once, since the browser asks for them together - for (const auto& other : textures) { - const auto otherBlock = other.starts_with('#') ? GeneralUtils::TryParse(other.substr(1)) : std::nullopt; - auto file = otherBlock ? NifFile::EmbeddedTexture(*data, *otherBlock) : std::nullopt; - if (!file) continue; - const auto bytes = file->size(); - embedded.Put(zone->assets[asset] + other, std::make_shared(std::move(*file)), bytes); - } - dds = NifFile::EmbeddedTexture(*data, *block); - } + void ReplyTexture(HTTPReply& reply, const HTTPContext& context, uint32_t zoneId, uint32_t asset, uint32_t slot, uint32_t lod) { + const auto* found = AssetPath(reply, zoneId, asset); + if (!found) return; + const auto path = *found; + lod = std::min(lod, MAX_LOD); + Touch(zoneId, lod); + const auto file = ClientAssets::ResolveResFile(path); + if (!file) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such texture"); + Files(); + // Quick when the model is converted already and the texture is a file of its own, or one kept from its model + auto priority = WorkerPool::ePriority::URGENT; + if (const auto cached = g_Models.Get(ModelKey(path, lod))) { + const auto textures = TexturesOf(*cached); + if (slot < textures.size() && textures[slot].starts_with('#') && !g_Embedded.Get(path + textures[slot])) priority = PriorityOf(*Zone(zoneId), path, *file); } else { - dds = ClientAssets::ReadResFile(texture); + priority = PriorityOf(*Zone(zoneId), path, *file); } - if (!dds) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "Could not read this texture"); - Binary(reply, std::move(*dds)); + const auto deferred = Web::Defer(reply, context); + g_Pool.Submit(priority, [deferred, path, lod, slot, file = *file, res = ClientAssets::ResFolder()] { + if (deferred.Cancelled()) return; + HTTPReply out; + const auto encoded = Encoded(path, lod, file); + const auto textures = encoded ? TexturesOf(*encoded) : std::vector{}; + if (slot >= textures.size() || textures[slot].empty()) { + JsonError(out, eHTTPStatusCode::NOT_FOUND, "No such texture"); + } else if (auto dds = TextureBytes(path, file, textures, textures[slot], res)) { + Binary(out, std::move(*dds)); + } else { + JsonError(out, eHTTPStatusCode::NOT_FOUND, "Could not read this texture"); + } + deferred.Send(std::move(out)); + }); + } + + void Shutdown() { + g_Pool.Stop(); } void RegisterRoutes() { + auto threads = Game::config ? Game::config->GetValue("scenery_workers", 0) : 0; + if (threads == 0) threads = static_cast(WorkerPool::DefaultThreads(std::thread::hardware_concurrency())); + threads = std::clamp(threads, 2, 16); + // Converting ahead of time never takes more than half of the threads besides the fast lane + g_Pool.Start(threads, std::max(1, (threads - 1) / 2)); + LOG("Converting scenery models with %u threads", threads); + Route(eHTTPMethod::GET, "/api/scenery/:zone/mesh/:asset", 0, "Model `asset` of a zone's scenery (see the scenery routes of properties and /world3d), converted from the client's .nif. Query: ?lod=0 (most detailed) to 3", [](HTTPReply& reply, const HTTPContext& context) { const auto zone = PathId(context.path, 2); const auto asset = PathId(context.path, 4); if (!zone || !asset) return JsonError(reply, eHTTPStatusCode::BAD_REQUEST, "Invalid zone or model"); - ReplyMesh(reply, *zone, *asset, LodOf(context)); + ReplyMesh(reply, context, *zone, *asset, LodOf(context)); }); Route(eHTTPMethod::GET, "/api/scenery/:zone/texture/:asset/:slot", 0, @@ -490,7 +753,7 @@ namespace Scenery { const auto asset = PathId(context.path, 4); const auto slot = PathId(context.path, 5); if (!zone || !asset || !slot) return JsonError(reply, eHTTPStatusCode::BAD_REQUEST, "Invalid zone, model or texture"); - ReplyTexture(reply, *zone, *asset, *slot, LodOf(context)); + ReplyTexture(reply, context, *zone, *asset, *slot, LodOf(context)); }); } } diff --git a/dDashboardServer/routes/Scenery.h b/dDashboardServer/routes/Scenery.h index 293612d33..4e3ae0b8c 100644 --- a/dDashboardServer/routes/Scenery.h +++ b/dDashboardServer/routes/Scenery.h @@ -5,6 +5,7 @@ #include struct HTTPReply; +struct HTTPContext; namespace WorldScene { struct Object; @@ -19,6 +20,12 @@ namespace WorldScene { * wants, nearest first. Models are converted from .nif on the server (NifFile) to a small binary format and kept in a * bounded in-memory cache; textures are sent as the client's own DDS files (textures stored inside a .nif are * wrapped as DDS), which the browser decodes itself, so nothing needs ImageMagick. + * + * Converting a big model takes a while, so models and textures not in memory are answered from a few worker threads + * (Web::Defer, WorkerPool; scenery_workers in dashboardconfig.ini) and the web server keeps answering meanwhile. + * Flairs and small models go first, in a lane of their own; the same model asked for twice at once is converted + * once. When a zone's manifest is asked for, its models are converted ahead of time onto the disk cache while + * someone views the zone. */ namespace Scenery { /** @@ -40,12 +47,17 @@ namespace Scenery { // Whether the client draws a model for this scene object (its render component's, or its nif_name) bool HasModel(const WorldScene::Object& object); - // Reply with model `asset` of the zone's manifest (NifFile::Encode), `lod` 0 the most detailed - void ReplyMesh(HTTPReply& reply, uint32_t zoneId, uint32_t asset, uint32_t lod); + // Reply with model `asset` of the zone's manifest (NifFile::Encode), `lod` 0 the most detailed. Deferred (answered + // from a worker thread) unless the model is in memory. + void ReplyMesh(HTTPReply& reply, const HTTPContext& context, uint32_t zoneId, uint32_t asset, uint32_t lod); - // Reply with texture `slot` (the model's "textures" list at that `lod`) of model `asset`, as a DDS file - void ReplyTexture(HTTPReply& reply, uint32_t zoneId, uint32_t asset, uint32_t slot, uint32_t lod); + // Reply with texture `slot` (the model's "textures" list at that `lod`) of model `asset`, as a DDS file. Deferred. + void ReplyTexture(HTTPReply& reply, const HTTPContext& context, uint32_t zoneId, uint32_t asset, uint32_t slot, uint32_t lod); - // /api/scenery/:zone/mesh/:asset and /api/scenery/:zone/texture/:asset/:slot, for anyone signed in + // /api/scenery/:zone/mesh/:asset and /api/scenery/:zone/texture/:asset/:slot, for anyone signed in; starts the + // conversion threads void RegisterRoutes(); + + // Stop the conversion threads (queued work is dropped) + void Shutdown(); } diff --git a/dDashboardServer/routes/SettingsCatalog.cpp b/dDashboardServer/routes/SettingsCatalog.cpp index 938e61b39..4b6b78a5a 100644 --- a/dDashboardServer/routes/SettingsCatalog.cpp +++ b/dDashboardServer/routes/SettingsCatalog.cpp @@ -278,6 +278,9 @@ namespace { c.Add(Bool(DASHBOARD, "secure_cookies", "HTTPS only cookies", "Turn on when the dashboard is served over HTTPS.", false)); c.Add(Bool(DASHBOARD, "behind_proxy", "Behind a reverse proxy", "Use the proxy's X-Forwarded-For for rate limits. Only when the dashboard can't be reached directly.", false)); c.Add(Unit(Int(DASHBOARD, "broadcast_interval", "Live update interval", "How often server status is pushed to open pages.", "2000", 250, 60000, true), "ms")); + c.Add(Unit(Int(DASHBOARD, "scenery_workers", "3D model conversion threads", + "Threads converting the client's models for the 3D views, so the dashboard keeps answering meanwhile. One of them only takes flairs and small models. 0 picks half the CPU cores (2 to 4).", + "0", 0, 16, true), "threads")); c.AddSection("Keys", "Only in the .ini file or the environment."); c.Add(Secret(DASHBOARD, "jwt_secret", "Session signing secret", "At least 32 characters; empty makes one. Changing it signs everyone out.", true)); diff --git a/dDashboardServer/routes/Showcase.cpp b/dDashboardServer/routes/Showcase.cpp index c38f0944c..312a83f8e 100644 --- a/dDashboardServer/routes/Showcase.cpp +++ b/dDashboardServer/routes/Showcase.cpp @@ -342,7 +342,7 @@ void RegisterShowcaseRoutes() { if (!zone) return; const auto asset = PathId(context.path, 5); if (!asset) return JsonError(reply, eHTTPStatusCode::BAD_REQUEST, "Invalid model"); - Scenery::ReplyMesh(reply, *zone, *asset, lodOf(context)); + Scenery::ReplyMesh(reply, context, *zone, *asset, lodOf(context)); }); Route(eHTTPMethod::GET, "/api/showcase/scenery/:zone/texture/:asset/:slot", PUBLIC, "A scenery texture of a property zone as DDS for the showcase's 3D view", @@ -352,6 +352,6 @@ void RegisterShowcaseRoutes() { const auto asset = PathId(context.path, 5); const auto slot = PathId(context.path, 6); if (!asset || !slot) return JsonError(reply, eHTTPStatusCode::BAD_REQUEST, "Invalid model or texture"); - Scenery::ReplyTexture(reply, *zone, *asset, *slot, lodOf(context)); + Scenery::ReplyTexture(reply, context, *zone, *asset, *slot, lodOf(context)); }); } diff --git a/dDashboardServer/routes/WorkerPool.cpp b/dDashboardServer/routes/WorkerPool.cpp new file mode 100644 index 000000000..6a6f43e4d --- /dev/null +++ b/dDashboardServer/routes/WorkerPool.cpp @@ -0,0 +1,140 @@ +#include "WorkerPool.h" + +#include +#include + +#include "Game.h" +#include "Logger.h" + +std::optional WorkerPool::Pick(const std::array& queued, bool fastLane, size_t runningBackground, size_t maxBackground) { + for (size_t i = 0; i < PRIORITIES; i++) { + const auto priority = static_cast(i); + if (fastLane && priority != ePriority::URGENT) break; + if (priority == ePriority::BACKGROUND && runningBackground >= maxBackground) break; + if (queued[i] > 0) return priority; + } + return std::nullopt; +} + +size_t WorkerPool::DefaultThreads(size_t cores) { + return std::clamp(cores / 2, 2, 4); +} + +void WorkerPool::Start(size_t threads, size_t maxBackground) { + Stop(); + std::lock_guard lock(m_Mutex); + m_Stopping = false; + m_MaxBackground = std::max(1, maxBackground); + threads = std::max(2, threads); + for (size_t i = 0; i < threads; i++) m_Threads.emplace_back([this, i] { Work(i == 0); }); +} + +void WorkerPool::Stop() { + std::vector threads; + { + std::lock_guard lock(m_Mutex); + m_Stopping = true; + for (auto& queue : m_Queues) queue.clear(); + threads.swap(m_Threads); + } + m_Wake.notify_all(); + for (auto& thread : threads) thread.join(); + m_Idle.notify_all(); +} + +void WorkerPool::Submit(ePriority priority, Job job, uint64_t group, bool front) { + if (!job) return; + { + std::lock_guard lock(m_Mutex); + if (!m_Threads.empty()) { + auto& queue = m_Queues[static_cast(priority)]; + Entry entry{ std::move(job), group }; + if (front) queue.push_front(std::move(entry)); + else queue.push_back(std::move(entry)); + job = nullptr; + } + } + if (!job) { + // Wake them all: the one woken might be the fast lane, which can't take this + m_Wake.notify_all(); + return; + } + try { + job(); + } catch (const std::exception& ex) { + LOG("A background job failed: %s", ex.what()); + } +} + +size_t WorkerPool::Cancel(uint64_t group) { + if (group == 0) return 0; + size_t dropped = 0; + { + std::lock_guard lock(m_Mutex); + for (auto& queue : m_Queues) { + const auto before = queue.size(); + std::erase_if(queue, [group](const Entry& entry) { return entry.group == group; }); + dropped += before - queue.size(); + } + } + m_Idle.notify_all(); + return dropped; +} + +size_t WorkerPool::Queued() const { + std::lock_guard lock(m_Mutex); + size_t total = 0; + for (const auto& queue : m_Queues) total += queue.size(); + return total; +} + +size_t WorkerPool::Active() const { + std::lock_guard lock(m_Mutex); + return m_Active; +} + +void WorkerPool::WaitIdle() { + std::unique_lock lock(m_Mutex); + m_Idle.wait(lock, [this] { + if (m_Active > 0) return false; + return std::all_of(m_Queues.begin(), m_Queues.end(), [](const auto& queue) { return queue.empty(); }); + }); +} + +void WorkerPool::Work(bool fastLane) { + std::unique_lock lock(m_Mutex); + while (true) { + std::optional priority; + m_Wake.wait(lock, [&] { + if (m_Stopping) return true; + std::array queued{}; + for (size_t i = 0; i < PRIORITIES; i++) queued[i] = m_Queues[i].size(); + priority = Pick(queued, fastLane, m_RunningBackground, m_MaxBackground); + return priority.has_value(); + }); + if (m_Stopping) return; + + auto& queue = m_Queues[static_cast(*priority)]; + auto entry = std::move(queue.front()); + queue.pop_front(); + const bool background = *priority == ePriority::BACKGROUND; + if (background) m_RunningBackground++; + m_Active++; + lock.unlock(); + + try { + entry.job(); + } catch (const std::exception& ex) { + LOG("A background job failed: %s", ex.what()); + } + entry.job = nullptr; // what it holds is released outside the lock + + lock.lock(); + m_Active--; + if (background) { + m_RunningBackground--; + m_Wake.notify_all(); // another background job may start now + } + m_Idle.notify_all(); + } +} diff --git a/dDashboardServer/routes/WorkerPool.h b/dDashboardServer/routes/WorkerPool.h new file mode 100644 index 000000000..7c9e46007 --- /dev/null +++ b/dDashboardServer/routes/WorkerPool.h @@ -0,0 +1,95 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +/** + * A few worker threads for slow work that routes hand off with Web::Defer (converting the client's models), so the + * web server's one thread keeps answering. Jobs run by priority, in the order they came within one. + * + * The first thread is a fast lane that only runs URGENT jobs, so something small and wanted now (a flair near the + * camera) never waits behind big conversions filling the other threads. BACKGROUND jobs (converting ahead of time) + * only run when nothing else waits, and only a few at once, so they never take every thread. + */ +class WorkerPool { +public: + enum class ePriority : uint8_t { + URGENT, // small and wanted now: the fast lane takes these too + NORMAL, + LARGE, // big jobs someone waits for + BACKGROUND, // nobody waits for it + }; + static constexpr size_t PRIORITIES = 4; + + using Job = std::function; + + WorkerPool() = default; + ~WorkerPool() { Stop(); } + WorkerPool(const WorkerPool&) = delete; + WorkerPool& operator=(const WorkerPool&) = delete; + + /** + * Start `threads` workers (at least 2: the fast lane and one more). At most `maxBackground` BACKGROUND jobs run + * at once (at least 1). + */ + void Start(size_t threads, size_t maxBackground = 1); + + // Drop the queued jobs and wait for the running ones + void Stop(); + + bool Running() const { return !m_Threads.empty(); } + size_t Threads() const { return m_Threads.size(); } + + /** + * Queue a job. `group` (0 for none) lets Cancel drop jobs queued together; `front` puts it before the others of + * its priority. Without threads (not started) the job runs right away on the caller's thread. + */ + void Submit(ePriority priority, Job job, uint64_t group = 0, bool front = false); + + // Drop the queued jobs of `group`; returns how many + size_t Cancel(uint64_t group); + + size_t Queued() const; + size_t Active() const; + + // Wait until nothing is queued or running (for tests) + void WaitIdle(); + + /** + * Which priority a worker takes next from queues of these lengths: the most urgent one waiting, only URGENT for + * the fast lane, and BACKGROUND only while fewer than maxBackground of those run. nullopt: nothing for it. + */ + static std::optional Pick(const std::array& queued, bool fastLane, size_t runningBackground, size_t maxBackground); + + /** + * The default number of threads for this many CPU cores: half of them, from 2 to 4 (a conversion is one core's + * work, and the game servers on the same machine need the rest) + */ + static size_t DefaultThreads(size_t cores); + +private: + struct Entry { + Job job; + uint64_t group{}; + }; + + void Work(bool fastLane); + + mutable std::mutex m_Mutex; + std::condition_variable m_Wake; + std::condition_variable m_Idle; + std::array, PRIORITIES> m_Queues; + std::vector m_Threads; + size_t m_MaxBackground{ 1 }; + size_t m_RunningBackground{}; + size_t m_Active{}; + bool m_Stopping{}; +}; diff --git a/dWeb/CMakeLists.txt b/dWeb/CMakeLists.txt index ae681a4ea..b0179321e 100644 --- a/dWeb/CMakeLists.txt +++ b/dWeb/CMakeLists.txt @@ -1,5 +1,6 @@ set(DWEB_SOURCES - "Web.cpp") + "Web.cpp" + "DeferredReply.cpp") add_library(dWeb STATIC ${DWEB_SOURCES}) diff --git a/dWeb/DeferredReply.cpp b/dWeb/DeferredReply.cpp new file mode 100644 index 000000000..59802a791 --- /dev/null +++ b/dWeb/DeferredReply.cpp @@ -0,0 +1,69 @@ +#include "DeferredReply.h" + +#include + +void DeferredReply::Send(HTTPReply reply) const { + if (!m_State || m_State->answered.exchange(true) || m_State->cancelled.load() || !m_State->queue) return; + reply.deferred.reset(); + m_State->queue->Push(*m_State, std::move(reply)); +} + +std::shared_ptr DeferredQueue::Begin(unsigned long connection) { + auto state = std::make_shared(); + state->connection = connection; + state->id = m_NextId++; + state->queue = this; + // A connection answers one request at a time, so an older entry is a request that is no longer waited on + if (const auto it = m_Pending.find(connection); it != m_Pending.end()) it->second.state->cancelled = true; + m_Pending[connection] = Waiting{ state, {}, false }; + return state; +} + +void DeferredQueue::SetReplyOptions(unsigned long connection, std::vector headers, bool close) { + const auto it = m_Pending.find(connection); + if (it == m_Pending.end()) return; + it->second.headers = std::move(headers); + it->second.close = close; +} + +void DeferredQueue::Close(unsigned long connection) { + const auto it = m_Pending.find(connection); + if (it == m_Pending.end()) return; + it->second.state->cancelled = true; + m_Pending.erase(it); +} + +void DeferredQueue::CancelAll() { + for (auto& [connection, waiting] : m_Pending) waiting.state->cancelled = true; + m_Pending.clear(); + std::lock_guard lock(m_Mutex); + m_Done.clear(); +} + +std::vector DeferredQueue::Drain() { + std::vector> done; + { + std::lock_guard lock(m_Mutex); + done.swap(m_Done); + } + std::vector finished; + if (done.empty()) return finished; + std::unordered_map connectionOf; + for (const auto& [connection, waiting] : m_Pending) connectionOf.emplace(waiting.state->id, connection); + for (auto& [id, reply] : done) { + const auto it = connectionOf.find(id); + if (it == connectionOf.end()) continue; // the client went away first + auto waiting = m_Pending.find(it->second); + auto headers = std::move(waiting->second.headers); + headers.insert(headers.end(), std::make_move_iterator(reply.headers.begin()), std::make_move_iterator(reply.headers.end())); + reply.headers = std::move(headers); + finished.push_back({ it->second, std::move(reply), waiting->second.close }); + m_Pending.erase(waiting); + } + return finished; +} + +void DeferredQueue::Push(const DeferredState& state, HTTPReply reply) { + std::lock_guard lock(m_Mutex); + m_Done.emplace_back(state.id, std::move(reply)); +} diff --git a/dWeb/DeferredReply.h b/dWeb/DeferredReply.h new file mode 100644 index 000000000..fb87defa8 --- /dev/null +++ b/dWeb/DeferredReply.h @@ -0,0 +1,93 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include + +#include "HTTPReply.h" + +class DeferredQueue; + +// Shared between the web thread and whoever answers a deferred request +struct DeferredState { + unsigned long connection{}; + uint64_t id{}; // unique per request, so a late answer never reaches a later request + DeferredQueue* queue{}; + std::atomic cancelled{}; + std::atomic answered{}; +}; + +/** + * A request a route answers later, from any thread: for slow work (converting a big model) that would otherwise hold + * up every other request, since the web server answers requests one at a time on one thread. Get one with + * Web::Defer, hand it to a worker, and Send the reply when the work is done; the web thread sends it on its next + * poll. When the client has gone first, Cancelled() turns true and the reply is dropped. + */ +class DeferredReply { +public: + DeferredReply() = default; + explicit DeferredReply(std::shared_ptr state) : m_State(std::move(state)) {} + + // Answer the request (any thread). Only the first answer counts; it is dropped when the client has gone. + void Send(HTTPReply reply) const; + + // Whether the client has gone (closed the connection, or the server is stopping), so the work can be skipped + bool Cancelled() const { return !m_State || m_State->cancelled.load(); } + + explicit operator bool() const { return m_State != nullptr; } + +private: + std::shared_ptr m_State; +}; + +/** + * Deferred requests waiting for their answers, and the answers that have arrived. Everything but Push (which + * DeferredReply::Send calls from any thread) is for the web thread. No sockets, so it can be unit tested. + */ +class DeferredQueue { +public: + struct Finished { + unsigned long connection{}; + HTTPReply reply; + bool close{}; // the client asked for Connection: close + }; + + // The request on `connection` will be answered later + std::shared_ptr Begin(unsigned long connection); + + // Headers the reply is sent with besides its own (e.g. a session cookie middleware refreshed), and whether the + // connection closes after it + void SetReplyOptions(unsigned long connection, std::vector headers, bool close); + + // The connection closed (or the request was answered another way): a late answer is dropped + void Close(unsigned long connection); + + // Cancel every pending request (the server is stopping) + void CancelAll(); + + // The answers that arrived for requests still pending; those requests stop being pending + std::vector Drain(); + + bool IsPending(unsigned long connection) const { return m_Pending.contains(connection); } + size_t Pending() const { return m_Pending.size(); } + + // DeferredReply::Send's half: any thread + void Push(const DeferredState& state, HTTPReply reply); + +private: + struct Waiting { + std::shared_ptr state; + std::vector headers; + bool close{}; + }; + + std::map m_Pending; // web thread only + uint64_t m_NextId{ 1 }; + + std::mutex m_Mutex; + std::vector> m_Done; // guarded by m_Mutex, keyed by DeferredState::id +}; diff --git a/dWeb/HTTPContext.h b/dWeb/HTTPContext.h index 819b7265a..5283fd5a7 100644 --- a/dWeb/HTTPContext.h +++ b/dWeb/HTTPContext.h @@ -27,6 +27,7 @@ struct HTTPContext { // Client information std::string clientIP{}; + unsigned long connectionId = 0; // the web server's id of the connection (for Web::Defer) // Authentication information (populated by auth middleware) bool isAuthenticated = false; diff --git a/dWeb/HTTPReply.h b/dWeb/HTTPReply.h new file mode 100644 index 000000000..1083122c7 --- /dev/null +++ b/dWeb/HTTPReply.h @@ -0,0 +1,37 @@ +#pragma once + +#include +#include +#include + +#include "eHTTPStatusCode.h" + +// Content type enum for HTTP responses +enum class eContentType { + APPLICATION_JSON, + TEXT_HTML, + TEXT_CSS, + TEXT_JAVASCRIPT, + TEXT_PLAIN, + TEXT_CSV, + IMAGE_PNG, + IMAGE_JPEG, + APPLICATION_OCTET_STREAM, + TEXT_PROMETHEUS // the Prometheus text exposition format +}; + +struct DeferredState; + +// For passing HTTP messages between functions +struct HTTPReply { + eHTTPStatusCode status = eHTTPStatusCode::NOT_FOUND; + std::string message = "{\"error\":\"Not Found\"}"; + eContentType contentType = eContentType::APPLICATION_JSON; + std::string location = ""; // For redirect responses (Location header) + std::vector headers{}; // Extra raw headers, e.g. "Set-Cookie: a=b" + // When set on a 200 reply, this file is streamed from disk as the body (with contentType and headers) instead of + // message, so large downloads never sit in memory + std::string file{}; + // Set by Web::Defer: the handler answers later, from another thread (DeferredReply), so nothing is sent now + std::shared_ptr deferred{}; +}; diff --git a/dWeb/Web.cpp b/dWeb/Web.cpp index 1c0b02059..9c0e19ca0 100644 --- a/dWeb/Web.cpp +++ b/dWeb/Web.cpp @@ -72,6 +72,9 @@ namespace { // Global middleware applied to all routes std::vector g_GlobalMiddleware; + + // Requests answered later by Web::Defer + DeferredQueue g_Deferred; // Helper to extract client IP from mongoose connection static std::string GetClientIP(mg_connection* connection) { @@ -124,6 +127,7 @@ namespace { // Get client IP context.clientIP = GetClientIP(connection); + context.connectionId = connection ? connection->id : 0; } const char* ContentTypeToString(eContentType contentType) { @@ -180,6 +184,53 @@ namespace { } } +// Send a reply; http_msg is the request (null for a deferred reply, whose request is gone) +static void SendReply(mg_connection* connection, const HTTPReply& reply, const mg_http_message* http_msg) { + // Build headers + std::string headers = std::string("Content-Type: ") + ContentTypeToString(reply.contentType) + "\r\n"; + if (!reply.location.empty()) { + headers += "Location: " + reply.location + "\r\n"; + } + // A route's own header replaces a default header of the same name + const auto headerName = [](const std::string& header) { + std::string name = header.substr(0, header.find(':')); + std::transform(name.begin(), name.end(), name.begin(), ::tolower); + return name; + }; + for (const auto& header : Game::web.GetDefaultHeaders()) { + const auto name = headerName(header); + if (std::ranges::none_of(reply.headers, [&](const std::string& h) { return headerName(h) == name; })) headers += header + "\r\n"; + } + for (const auto& header : reply.headers) headers += header + "\r\n"; + + if (!reply.file.empty() && reply.status == eHTTPStatusCode::OK) { + // Streamed in chunks by mongoose. Content-Type comes from the mime override (it adds its own header). + std::string extraHeaders = headers.substr(headers.find("\r\n") + 2); + const std::string mimeTypes = std::string("*=") + ContentTypeToString(reply.contentType); + mg_http_serve_opts opts{}; + opts.extra_headers = extraHeaders.c_str(); + opts.mime_types = mimeTypes.c_str(); + // Without Accept-Encoding, so a stale ".gz" next to the file is never served instead + mg_http_message request{}; + if (http_msg) request = *http_msg; + else request.method = mg_str("GET"); + for (auto& header : request.headers) { + if (header.name.len && mg_strcasecmp(header.name, mg_str("Accept-Encoding")) == 0) header.name = mg_str("X-Ignored"); + } + mg_http_serve_file(connection, &request, reply.file.c_str(), &opts); + return; + } + + // Written by hand rather than with mg_http_reply: that pads Content-Length with spaces (it fills the number in + // afterwards), which strict clients such as Node's fetch reject, and it can't send binary bodies + headers += "Content-Length: " + std::to_string(reply.message.size()) + "\r\n"; + const auto status = static_cast(reply.status); + std::string resp = "HTTP/1.1 " + std::to_string(status) + " " + ReasonPhrase(status) + "\r\n" + headers + "\r\n"; + mg_send(connection, resp.data(), resp.size()); + mg_send(connection, reply.message.data(), reply.message.size()); + connection->is_resp = 0; +} + void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_msg) { if (g_HTTPRoutes.empty()) return; @@ -394,47 +445,16 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms } } - // Build headers - std::string headers = std::string("Content-Type: ") + ContentTypeToString(reply.contentType) + "\r\n"; - if (!reply.location.empty()) { - headers += "Location: " + reply.location + "\r\n"; - } - // A route's own header replaces a default header of the same name - const auto headerName = [](const std::string& header) { - std::string name = header.substr(0, header.find(':')); - std::transform(name.begin(), name.end(), name.begin(), ::tolower); - return name; - }; - for (const auto& header : Game::web.GetDefaultHeaders()) { - const auto name = headerName(header); - if (std::ranges::none_of(reply.headers, [&](const std::string& h) { return headerName(h) == name; })) headers += header + "\r\n"; - } - for (const auto& header : reply.headers) headers += header + "\r\n"; - - if (!reply.file.empty() && reply.status == eHTTPStatusCode::OK && http_msg) { - // Streamed in chunks by mongoose. Content-Type comes from the mime override (it adds its own header). - std::string extraHeaders = headers.substr(headers.find("\r\n") + 2); - const std::string mimeTypes = std::string("*=") + ContentTypeToString(reply.contentType); - mg_http_serve_opts opts{}; - opts.extra_headers = extraHeaders.c_str(); - opts.mime_types = mimeTypes.c_str(); - // Without Accept-Encoding, so a stale ".gz" next to the file is never served instead - mg_http_message request = *http_msg; - for (auto& header : request.headers) { - if (header.name.len && mg_strcasecmp(header.name, mg_str("Accept-Encoding")) == 0) header.name = mg_str("X-Ignored"); - } - mg_http_serve_file(connection, &request, reply.file.c_str(), &opts); + if (reply.deferred) { + // Answered later (Web::Defer): the connection keeps is_resp set, so mongoose reads no further request on it + const auto* cc = http_msg ? mg_http_get_header(const_cast(http_msg), "Connection") : nullptr; + g_Deferred.SetReplyOptions(connection->id, reply.headers, cc && mg_strcasecmp(*cc, mg_str("close")) == 0); return; } + // The handler deferred and then failed: its late answer is dropped + if (g_Deferred.IsPending(connection->id)) g_Deferred.Close(connection->id); - // Written by hand rather than with mg_http_reply: that pads Content-Length with spaces (it fills the number in - // afterwards), which strict clients such as Node's fetch reject, and it can't send binary bodies - headers += "Content-Length: " + std::to_string(reply.message.size()) + "\r\n"; - const auto status = static_cast(reply.status); - std::string resp = "HTTP/1.1 " + std::to_string(status) + " " + ReasonPhrase(status) + "\r\n" + headers + "\r\n"; - mg_send(connection, resp.data(), resp.size()); - mg_send(connection, reply.message.data(), reply.message.size()); - connection->is_resp = 0; + SendReply(connection, reply, http_msg); } @@ -554,6 +574,8 @@ void HandleMessages(mg_connection* connection, int message, void* message_data) break; case MG_EV_CLOSE: g_AuthenticatedWSConnections.erase(connection); + // A deferred answer that comes after this is dropped + g_Deferred.Close(connection->id); break; default: break; @@ -638,7 +660,21 @@ Web::Web() { } Web::~Web() { + // Static destruction: the maps and queues the close events touch are gone by now, so the handlers are off + // (HandleMessages checks enabled). Servers call Shutdown first, while they're still there. + enabled = false; + if (!managerFreed) mg_mgr_free(&mgr); + managerFreed = true; +} + +void Web::Shutdown() { + if (managerFreed) return; + // Closing the connections fires their close events, which still clean up (WebSocket clients, deferred requests) + g_Deferred.CancelAll(); mg_mgr_free(&mgr); + managerFreed = true; + enabled = false; + g_AuthenticatedWSConnections.clear(); } bool Web::Startup(const std::string& listen_ip, const uint32_t listen_port) { @@ -677,9 +713,30 @@ bool Web::Startup(const std::string& listen_ip, const uint32_t listen_port) { void Web::ReceiveRequests(int timeoutMs) { mg_mgr_poll(&mgr, timeoutMs); + SendDeferredReplies(); RecheckDueWebSockets(); } +void Web::SendDeferredReplies() { + for (auto& finished : g_Deferred.Drain()) { + mg_connection* connection = mgr.conns; + while (connection && connection->id != finished.connection) connection = connection->next; + if (!connection || connection->is_closing) continue; + // Clears is_resp once the reply is out, so mongoose reads the connection's next request again + SendReply(connection, finished.reply, nullptr); + if (finished.close) connection->is_draining = 1; + } +} + +DeferredReply Web::Defer(HTTPReply& reply, const HTTPContext& context) { + reply.deferred = g_Deferred.Begin(context.connectionId); + return DeferredReply(reply.deferred); +} + +size_t Web::PendingDeferred() const { + return g_Deferred.Pending(); +} + void Web::RecheckWebSockets(uint32_t accountId) { for (auto& [connection, client] : g_AuthenticatedWSConnections) { if (accountId == 0 || client.accountId == accountId) client.nextCheck = {}; diff --git a/dWeb/Web.h b/dWeb/Web.h index 4f4377020..5a7a1baec 100644 --- a/dWeb/Web.h +++ b/dWeb/Web.h @@ -10,6 +10,8 @@ #include "json_fwd.hpp" #include "eHTTPStatusCode.h" #include "HTTPContext.h" +#include "HTTPReply.h" +#include "DeferredReply.h" #include "IHTTPMiddleware.h" // Forward declarations for game namespace @@ -24,32 +26,6 @@ enum class eHTTPMethod; // Forward declaration for mongoose manager typedef struct mg_mgr mg_mgr; -// Content type enum for HTTP responses -enum class eContentType { - APPLICATION_JSON, - TEXT_HTML, - TEXT_CSS, - TEXT_JAVASCRIPT, - TEXT_PLAIN, - TEXT_CSV, - IMAGE_PNG, - IMAGE_JPEG, - APPLICATION_OCTET_STREAM, - TEXT_PROMETHEUS // the Prometheus text exposition format -}; - -// For passing HTTP messages between functions -struct HTTPReply { - eHTTPStatusCode status = eHTTPStatusCode::NOT_FOUND; - std::string message = "{\"error\":\"Not Found\"}"; - eContentType contentType = eContentType::APPLICATION_JSON; - std::string location = ""; // For redirect responses (Location header) - std::vector headers{}; // Extra raw headers, e.g. "Set-Cookie: a=b" - // When set on a 200 reply, this file is streamed from disk as the body (with contentType and headers) instead of - // message, so large downloads never sit in memory - std::string file{}; -}; - // HTTP route structure // This structure is used to register HTTP routes // with the server. Each route has a path, method, optional middleware, @@ -109,6 +85,14 @@ public: void RegisterWSSubscription(const std::string& subscription, uint8_t minLevel = 0); // The level is looked up each time (for permissions that can change while running) void RegisterWSSubscription(const std::string& subscription, std::function minLevel); + /** + * Answer this request later (from any thread) instead of when the handler returns, for slow work that would hold + * up every other request: the handler hands the returned DeferredReply to a worker and returns; the web thread + * sends the worker's reply on its next poll. Headers middleware set on `reply` are sent with it. + */ + static DeferredReply Defer(HTTPReply& reply, const HTTPContext& context); + // Deferred requests still waiting for their answers + size_t PendingDeferred() const; // Add global middleware that applies to all routes void AddGlobalMiddleware(MiddlewarePtr middleware); // Set WebSocket authentication callback for token validation @@ -117,6 +101,12 @@ public: void SetWSApiAccessCallback(std::function callback) { wsApiAccessCallback = std::move(callback); } // Returns if the web server is enabled bool IsEnabled() const { return enabled; }; + /** + * Close every connection and stop listening, before the server's other state is torn down (the destructor runs + * during static destruction, when what the connections' close events touch may be gone). Deferred requests still + * waiting are cancelled; stop the workers that answer them first. + */ + void Shutdown(); // Send a message to all connected WebSocket clients that are subscribed to the given topic void static SendWSMessage(std::string sub, nlohmann::json& message); // Send a message on a topic only to the subscribed connections of one account @@ -133,10 +123,14 @@ public: WSAuthCallback GetWSAuthCallback() const { return wsAuthCallback; } const std::function& GetWSApiAccessCallback() const { return wsApiAccessCallback; } private: + // Send the answers of deferred requests that have arrived + void SendDeferredReplies(); // mongoose manager mg_mgr mgr; // If the web server is enabled bool enabled = false; + // mg_mgr_free has run (Shutdown) + bool managerFreed = false; // WebSocket authentication callback WSAuthCallback wsAuthCallback = nullptr; std::function wsApiAccessCallback = nullptr; diff --git a/docs/Dashboard.md b/docs/Dashboard.md index c6068221d..f68fbb37f 100644 --- a/docs/Dashboard.md +++ b/docs/Dashboard.md @@ -894,8 +894,19 @@ Shadows and Hidden objects (off by default), remembered per account. scene object's model and the zone's sky, from the game client's files (needs `client_location`). Models load nearest first; the detail setting picks the model level of detail, how far objects are drawn (1400/800/450 units), texture sharpness and a memory budget. Switch it off with *Scenery*. The 3D world view draws the same as its *Models* layer -(on by default, Medium detail), plus the terrain's flairs. Converted models are cached in memory and in -`dDashboardServer/scenery_cache` next to the server (at most 512 MB); nothing needs ImageMagick. Endpoints: +(on by default, Medium detail), plus the terrain's flairs. Converted models are cached in memory (64 MB) and in +`dDashboardServer/scenery_cache` next to the server (at most 512 MB); nothing needs ImageMagick. + +Converting a model the caches don't have yet (a big "glom" file takes a moment) happens on a few worker threads, so +the dashboard keeps answering everything else meanwhile: the route hands the request to a worker (`Web::Defer`) and +the web thread sends the answer when it is ready. Flairs and small models (up to 256 KB) go first, and one of the +threads only takes those, so the grass around the camera never waits behind a big model; models of 4 MB and more wait +behind smaller ones. A model asked for twice at once (two viewers) is converted once. When a zone's scenery or flair +manifest is asked for, the zone's models are also converted ahead of time onto the disk cache, flairs and smallest +first, at the detail its viewer last used: only when nothing else waits, on at most half of the threads besides the +flairs' one, stopping when nobody has viewed the zone for 90 seconds or the disk cache is three quarters full (it +never evicts for this). The number of threads is `scenery_workers` in `dashboardconfig.ini` (Settings > Dashboard > +Web server; 0, the default, picks half the CPU cores, 2 to 4; read at startup). Endpoints: `/api/properties/:id/scenery`, `/api/world3d/:zone/scenery`, `/api/world3d/:zone/flairs`, `/api/scenery/:zone/mesh/:asset?lod=`, `/api/scenery/:zone/texture/:asset/:slot?lod=`. The world view's other data: `/api/world3d/:zone/scene` (objects and scenes), `/terrain_chunks` and `/terrain_layers` (the terrain file; sent diff --git a/tests/dWebTests/CMakeLists.txt b/tests/dWebTests/CMakeLists.txt index 074d8d55b..dbeb1546e 100644 --- a/tests/dWebTests/CMakeLists.txt +++ b/tests/dWebTests/CMakeLists.txt @@ -34,6 +34,9 @@ set(DWEBTESTS_SOURCES "${PROJECT_SOURCE_DIR}/dDashboardServer/ai/ModeratorPrompt.cpp" "WorldViewTests.cpp" "NifFileTests.cpp" + "DeferredReplyTests.cpp" + "WorkerPoolTests.cpp" + "${PROJECT_SOURCE_DIR}/dDashboardServer/routes/WorkerPool.cpp" "${PROJECT_SOURCE_DIR}/dDashboardServer/routes/NifFile.cpp" "BackupFilesTests.cpp" "SecurityFixesTests.cpp" diff --git a/tests/dWebTests/DeferredReplyTests.cpp b/tests/dWebTests/DeferredReplyTests.cpp new file mode 100644 index 000000000..a046cd357 --- /dev/null +++ b/tests/dWebTests/DeferredReplyTests.cpp @@ -0,0 +1,194 @@ +#include + +#include +#include +#include +#include +#include +#include + +#include "DeferredReply.h" +#include "Web.h" +#include "eHTTPMethod.h" + +namespace { + HTTPReply Text(const std::string& body) { + HTTPReply reply; + reply.status = eHTTPStatusCode::OK; + reply.contentType = eContentType::TEXT_PLAIN; + reply.message = body; + return reply; + } +} + +TEST(DeferredQueueTests, AnswerReachesItsConnection) { + DeferredQueue queue; + const DeferredReply deferred(queue.Begin(7)); + queue.SetReplyOptions(7, { "Set-Cookie: a=b" }, true); + EXPECT_TRUE(queue.Drain().empty()); + EXPECT_EQ(queue.Pending(), 1u); + + auto reply = Text("done"); + reply.headers.push_back("Cache-Control: no-store"); + deferred.Send(reply); + const auto finished = queue.Drain(); + ASSERT_EQ(finished.size(), 1u); + EXPECT_EQ(finished[0].connection, 7u); + EXPECT_EQ(finished[0].reply.message, "done"); + EXPECT_TRUE(finished[0].close); + // Middleware's headers first, then the handler's own + ASSERT_EQ(finished[0].reply.headers.size(), 2u); + EXPECT_EQ(finished[0].reply.headers[0], "Set-Cookie: a=b"); + EXPECT_EQ(finished[0].reply.headers[1], "Cache-Control: no-store"); + EXPECT_EQ(queue.Pending(), 0u); +} + +TEST(DeferredQueueTests, OnlyTheFirstAnswerCounts) { + DeferredQueue queue; + const DeferredReply deferred(queue.Begin(1)); + deferred.Send(Text("first")); + deferred.Send(Text("second")); + const auto finished = queue.Drain(); + ASSERT_EQ(finished.size(), 1u); + EXPECT_EQ(finished[0].reply.message, "first"); + EXPECT_TRUE(queue.Drain().empty()); +} + +TEST(DeferredQueueTests, ClosedConnectionCancelsAndDropsTheAnswer) { + DeferredQueue queue; + const DeferredReply deferred(queue.Begin(3)); + EXPECT_FALSE(deferred.Cancelled()); + queue.Close(3); + EXPECT_TRUE(deferred.Cancelled()); + deferred.Send(Text("late")); + EXPECT_TRUE(queue.Drain().empty()); + EXPECT_EQ(queue.Pending(), 0u); +} + +TEST(DeferredQueueTests, LateAnswerNeverReachesALaterRequest) { + DeferredQueue queue; + const DeferredReply first(queue.Begin(5)); + // The same connection starts another deferred request (the first was given up on) + const DeferredReply second(queue.Begin(5)); + EXPECT_TRUE(first.Cancelled()); + // Pushed straight to the queue, as a worker that checked before the cancel would + queue.Push(DeferredState{ .connection = 5, .id = 1 }, Text("stale")); + second.Send(Text("fresh")); + const auto finished = queue.Drain(); + ASSERT_EQ(finished.size(), 1u); + EXPECT_EQ(finished[0].reply.message, "fresh"); +} + +TEST(DeferredQueueTests, CancelAllCancelsEverything) { + DeferredQueue queue; + const DeferredReply a(queue.Begin(1)), b(queue.Begin(2)); + queue.CancelAll(); + EXPECT_TRUE(a.Cancelled()); + EXPECT_TRUE(b.Cancelled()); + EXPECT_EQ(queue.Pending(), 0u); +} + +TEST(DeferredQueueTests, AnswersFromManyThreads) { + DeferredQueue queue; + constexpr unsigned long COUNT = 64; + std::vector replies; + for (unsigned long i = 1; i <= COUNT; i++) replies.emplace_back(queue.Begin(i)); + std::vector threads; + for (unsigned long i = 0; i < COUNT; i++) threads.emplace_back([&replies, i] { replies[i].Send(Text(std::to_string(i + 1))); }); + for (auto& thread : threads) thread.join(); + const auto finished = queue.Drain(); + ASSERT_EQ(finished.size(), COUNT); + for (const auto& f : finished) EXPECT_EQ(f.reply.message, std::to_string(f.connection)); +} + +// The real web server: a deferred request doesn't hold up others, its answer arrives later on the same +// connection (which then takes its next request), and a client that leaves first is handled safely. +namespace { + struct Client { + std::vector bodies; + bool closed{}; + }; + + void ClientEvents(mg_connection* connection, int event, void* data) { + auto* client = static_cast(connection->fn_data); + if (event == MG_EV_HTTP_MSG) { + const auto* message = static_cast(data); + client->bodies.emplace_back(message->body.buf, message->body.len); + } else if (event == MG_EV_CLOSE) { + client->closed = true; + } + } + + void Get(mg_connection* connection, const std::string& path) { + mg_printf(connection, "GET %s HTTP/1.1\r\nHost: localhost\r\n\r\n", path.c_str()); + } +} + +TEST(DeferredWebTests, DeferredRequestsDontHoldUpOthers) { + constexpr uint32_t PORT = 38631; + ASSERT_TRUE(Game::web.Startup("127.0.0.1", PORT)); + + std::atomic release{}; + std::vector workers; + std::vector handedOff; + Game::web.RegisterHTTPRoute({ .path = "/slow", .method = eHTTPMethod::GET, .middleware = {}, .handle = [&](HTTPReply& reply, const HTTPContext& context) { + auto deferred = Web::Defer(reply, context); + handedOff.push_back(deferred); + workers.emplace_back([deferred, &release] { + while (!release) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + deferred.Send(Text("slow")); + }); + } }); + Game::web.RegisterHTTPRoute({ .path = "/fast", .method = eHTTPMethod::GET, .middleware = {}, .handle = [](HTTPReply& reply, const HTTPContext&) { + reply = Text("fast"); + } }); + + mg_mgr clients; + mg_mgr_init(&clients); + const auto poll = [&](const std::function& done) { + const auto until = std::chrono::steady_clock::now() + std::chrono::seconds(5); + while (!done() && std::chrono::steady_clock::now() < until) { + Game::web.ReceiveRequests(1); + mg_mgr_poll(&clients, 1); + } + return done(); + }; + const std::string url = "http://127.0.0.1:" + std::to_string(PORT); + + Client slow, fast, leaving; + auto* slowConnection = mg_http_connect(&clients, url.c_str(), ClientEvents, &slow); + auto* fastConnection = mg_http_connect(&clients, url.c_str(), ClientEvents, &fast); + auto* leavingConnection = mg_http_connect(&clients, url.c_str(), ClientEvents, &leaving); + ASSERT_NE(slowConnection, nullptr); + ASSERT_NE(fastConnection, nullptr); + ASSERT_NE(leavingConnection, nullptr); + Get(slowConnection, "/slow"); + Get(leavingConnection, "/slow"); + ASSERT_TRUE(poll([&] { return handedOff.size() == 2; })); + + // Answered while the slow ones still wait + Get(fastConnection, "/fast"); + ASSERT_TRUE(poll([&] { return !fast.bodies.empty(); })); + EXPECT_EQ(fast.bodies[0], "fast"); + EXPECT_TRUE(slow.bodies.empty()); + EXPECT_EQ(Game::web.PendingDeferred(), 2u); + + // One client leaves before its answer: the work sees it's cancelled, and the answer is dropped + leavingConnection->is_closing = 1; + ASSERT_TRUE(poll([&] { return handedOff[0].Cancelled() || handedOff[1].Cancelled(); })); + EXPECT_EQ(Game::web.PendingDeferred(), 1u); + + release = true; + ASSERT_TRUE(poll([&] { return !slow.bodies.empty(); })); + EXPECT_EQ(slow.bodies[0], "slow"); + EXPECT_EQ(Game::web.PendingDeferred(), 0u); + EXPECT_TRUE(leaving.bodies.empty()); + + // The connection takes its next request after the deferred answer + Get(slowConnection, "/fast"); + ASSERT_TRUE(poll([&] { return slow.bodies.size() == 2; })); + EXPECT_EQ(slow.bodies[1], "fast"); + + for (auto& worker : workers) worker.join(); + mg_mgr_free(&clients); +} diff --git a/tests/dWebTests/WorkerPoolTests.cpp b/tests/dWebTests/WorkerPoolTests.cpp new file mode 100644 index 000000000..b03a2accf --- /dev/null +++ b/tests/dWebTests/WorkerPoolTests.cpp @@ -0,0 +1,134 @@ +#include + +#include +#include +#include +#include +#include +#include + +#include "WorkerPool.h" + +using ePriority = WorkerPool::ePriority; + +TEST(WorkerPoolTests, PickTakesTheMostUrgent) { + EXPECT_EQ(WorkerPool::Pick({ 0, 1, 1, 1 }, false, 0, 1), ePriority::NORMAL); + EXPECT_EQ(WorkerPool::Pick({ 2, 1, 1, 1 }, false, 0, 1), ePriority::URGENT); + EXPECT_EQ(WorkerPool::Pick({ 0, 0, 3, 1 }, false, 0, 1), ePriority::LARGE); + EXPECT_EQ(WorkerPool::Pick({ 0, 0, 0, 1 }, false, 0, 1), ePriority::BACKGROUND); + EXPECT_EQ(WorkerPool::Pick({ 0, 0, 0, 0 }, false, 0, 1), std::nullopt); +} + +TEST(WorkerPoolTests, FastLaneOnlyTakesUrgentWork) { + EXPECT_EQ(WorkerPool::Pick({ 1, 1, 1, 1 }, true, 0, 1), ePriority::URGENT); + EXPECT_EQ(WorkerPool::Pick({ 0, 1, 1, 1 }, true, 0, 1), std::nullopt); +} + +TEST(WorkerPoolTests, BackgroundWorkIsLimited) { + EXPECT_EQ(WorkerPool::Pick({ 0, 0, 0, 5 }, false, 1, 1), std::nullopt); + EXPECT_EQ(WorkerPool::Pick({ 0, 0, 0, 5 }, false, 1, 2), ePriority::BACKGROUND); + // Other work still goes ahead + EXPECT_EQ(WorkerPool::Pick({ 0, 1, 0, 5 }, false, 1, 1), ePriority::NORMAL); +} + +TEST(WorkerPoolTests, DefaultThreads) { + EXPECT_EQ(WorkerPool::DefaultThreads(0), 2u); + EXPECT_EQ(WorkerPool::DefaultThreads(2), 2u); + EXPECT_EQ(WorkerPool::DefaultThreads(6), 3u); + EXPECT_EQ(WorkerPool::DefaultThreads(32), 4u); +} + +TEST(WorkerPoolTests, WithoutThreadsJobsRunRightAway) { + WorkerPool pool; + bool ran = false; + pool.Submit(ePriority::NORMAL, [&ran] { ran = true; }); + EXPECT_TRUE(ran); +} + +TEST(WorkerPoolTests, RunsEveryJob) { + WorkerPool pool; + pool.Start(3); + std::atomic count{}; + for (int i = 0; i < 200; i++) pool.Submit(static_cast(i % 4), [&count] { count++; }); + pool.WaitIdle(); + EXPECT_EQ(count.load(), 200); + EXPECT_EQ(pool.Queued(), 0u); +} + +// With every general thread busy on a big job, an urgent one still runs in the fast lane +TEST(WorkerPoolTests, UrgentWorkDoesNotWaitBehindBigJobs) { + WorkerPool pool; + pool.Start(2); + std::atomic release{}, bigStarted{}, urgentDone{}; + pool.Submit(ePriority::LARGE, [&] { + bigStarted = true; + while (!release) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + }); + while (!bigStarted) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + pool.Submit(ePriority::URGENT, [&] { urgentDone = true; }); + for (int i = 0; i < 2000 && !urgentDone; i++) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + EXPECT_TRUE(urgentDone.load()); + release = true; + pool.WaitIdle(); +} + +TEST(WorkerPoolTests, OrderWithinAndAcrossPriorities) { + WorkerPool pool; + pool.Start(2); + std::mutex mutex; + std::vector order; + std::atomic release{}, blockerStarted{}; + // Occupy the general thread so the rest queue up + pool.Submit(ePriority::NORMAL, [&] { + blockerStarted = true; + while (!release) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + }); + while (!blockerStarted) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + const auto job = [&](std::string name) { return [&, name] { std::lock_guard lock(mutex); order.push_back(name); }; }; + pool.Submit(ePriority::BACKGROUND, job("background")); + pool.Submit(ePriority::LARGE, job("large")); + pool.Submit(ePriority::NORMAL, job("normal 1")); + pool.Submit(ePriority::NORMAL, job("normal 2")); + pool.Submit(ePriority::NORMAL, job("normal 0"), 0, true); + release = true; + pool.WaitIdle(); + EXPECT_EQ(order, (std::vector{ "normal 0", "normal 1", "normal 2", "large", "background" })); +} + +TEST(WorkerPoolTests, CancelDropsAGroupsQueuedJobs) { + WorkerPool pool; + pool.Start(2); + std::atomic release{}, blockerStarted{}; + std::atomic ran{}; + pool.Submit(ePriority::NORMAL, [&] { + blockerStarted = true; + while (!release) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + }); + while (!blockerStarted) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + for (int i = 0; i < 5; i++) pool.Submit(ePriority::BACKGROUND, [&ran] { ran++; }, 42); + pool.Submit(ePriority::BACKGROUND, [&ran] { ran += 100; }, 7); + EXPECT_EQ(pool.Cancel(42), 5u); + EXPECT_EQ(pool.Cancel(0), 0u); + release = true; + pool.WaitIdle(); + EXPECT_EQ(ran.load(), 100); +} + +TEST(WorkerPoolTests, StopDropsQueuedJobs) { + WorkerPool pool; + pool.Start(2); + std::atomic release{}, blockerStarted{}; + std::atomic ran{}; + pool.Submit(ePriority::NORMAL, [&] { + blockerStarted = true; + while (!release) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + }); + while (!blockerStarted) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + for (int i = 0; i < 5; i++) pool.Submit(ePriority::NORMAL, [&ran] { ran++; }); + std::thread stopper([&pool] { pool.Stop(); }); + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + release = true; + stopper.join(); + EXPECT_EQ(ran.load(), 0); + EXPECT_FALSE(pool.Running()); +}