mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-03 11:23:43 +00:00
feat: convert 3D scenery models on worker threads with deferred replies
The dashboard's web server answers one request at a time, so converting a big .nif (glom files up to tens of MB) held up every other request, flairs included. - dWeb: Web::Defer hands a request to another thread; the reply is sent from the web thread on its next poll (DeferredQueue). A client that leaves first cancels it and the late reply is dropped. The synchronous route API is unchanged. - Web::Shutdown closes connections while the state their close events touch is still alive; the destructor no longer runs handlers during static destruction (stopping the dashboard aborted in ~WSClient). - WorkerPool: priority lanes, with one thread only for urgent work (flairs, small models, textures), and limited background work. - Scenery: mesh and texture routes (and the showcase's) convert on the pool; thread-safe memory and disk caches, one conversion per model at a time with waiters sharing it; zones are converted ahead onto the disk cache while viewed. - Setting scenery_workers (0: half the cores, 2 to 4). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -42,6 +42,7 @@ set(DASHBOARDROUTES_SOURCES
|
||||
"WorldView.cpp"
|
||||
"NifFile.cpp"
|
||||
"Scenery.cpp"
|
||||
"WorkerPool.cpp"
|
||||
"AccountModeration.cpp"
|
||||
"ModerationTools.cpp"
|
||||
"ModeratorHelper.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<std::string> 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<std::filesystem::path> 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<std::string> FindResFile(const std::string& folder, const std::string& fileName) {
|
||||
const auto normalized = NormalizeAssetPath(folder);
|
||||
if (ClientRes().empty() || !IsSafeAssetPath(normalized)) return std::nullopt;
|
||||
|
||||
@@ -37,6 +37,12 @@ namespace ClientAssets {
|
||||
// Where a file inside res/ is on disk (matched ignoring case), if it exists
|
||||
std::optional<std::filesystem::path> 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<std::filesystem::path> ResolveResFile(const std::string& relativePath, const std::filesystem::path& res);
|
||||
|
||||
// The res/-relative path of a file called `fileName` (ignoring case) anywhere under res/<folder>
|
||||
std::optional<std::string> FindResFile(const std::string& folder, const std::string& fileName);
|
||||
|
||||
|
||||
@@ -5,9 +5,13 @@
|
||||
#include <cstring>
|
||||
#include <filesystem>
|
||||
#include <fstream>
|
||||
#include <future>
|
||||
#include <map>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <set>
|
||||
#include <sstream>
|
||||
#include <thread>
|
||||
#include <unordered_map>
|
||||
#include <unordered_set>
|
||||
|
||||
@@ -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<std::string, size_t> index; // res path -> its index in assets
|
||||
std::string json;
|
||||
std::optional<std::string> flairs; // the flairs' manifest, built when first asked for (their models join assets)
|
||||
std::unordered_set<std::string> 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<uintmax_t> 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<std::pair<std::filesystem::file_time_type, std::filesystem::path>> 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<std::streamsize>(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<std::pair<std::filesystem::file_time_type, std::filesystem::path>> 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<std::streamsize>(data.size()))) return;
|
||||
|
||||
std::mutex m_Mutex;
|
||||
uintmax_t m_Total{};
|
||||
bool m_Counted{};
|
||||
};
|
||||
|
||||
DiskCache g_Disk;
|
||||
|
||||
using Bytes = std::shared_ptr<const std::string>;
|
||||
|
||||
// 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<std::string, Bytes> 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<std::string, std::shared_future<Bytes>> 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<const std::string>(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<std::string> textures; // per mesh: a res path, "#<block>" 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<const std::string>(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<const std::string> Encoded(const std::string& path, uint32_t lod) {
|
||||
static TtlCache<std::string, std::shared_ptr<const std::string>> 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<Bytes> promise;
|
||||
std::shared_future<Bytes> 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<std::string> textures; // per mesh: a res path, "#<block>" 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<const std::string>(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<uint32_t>(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<uint32_t, Activity> g_Activity;
|
||||
|
||||
void Touch(uint32_t zoneId, std::optional<uint32_t> 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<Activity> 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<uint64_t>(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<std::string>& 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<Item> items;
|
||||
std::set<std::string> 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<std::string> TextureBytes(const std::string& path, const std::filesystem::path& modelFile, const std::vector<std::string>& 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<int32_t>(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<int32_t>(other.substr(1)) : std::nullopt;
|
||||
auto file = otherBlock ? NifFile::EmbeddedTexture(*data, *otherBlock) : std::nullopt;
|
||||
if (file) g_Embedded.Put(path + other, std::make_shared<const std::string>(std::move(*file)));
|
||||
}
|
||||
auto dds = NifFile::EmbeddedTexture(*data, *block);
|
||||
if (dds) g_Embedded.Put(path + name, std::make_shared<const std::string>(*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<std::string> 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<std::string> 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<std::string>{};
|
||||
if (slot >= textures.size() || textures[slot].empty()) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such texture");
|
||||
const auto& texture = textures[slot];
|
||||
std::optional<std::string> 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<std::string, std::shared_ptr<const std::string>> 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<int32_t>(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<int32_t>(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<const std::string>(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<std::string>{};
|
||||
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<uint32_t>("scenery_workers", 0) : 0;
|
||||
if (threads == 0) threads = static_cast<uint32_t>(WorkerPool::DefaultThreads(std::thread::hardware_concurrency()));
|
||||
threads = std::clamp<uint32_t>(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<size_t>(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<uint32_t>(context.path, 2);
|
||||
const auto asset = PathId<uint32_t>(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<uint32_t>(context.path, 4);
|
||||
const auto slot = PathId<uint32_t>(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));
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
#include <string>
|
||||
|
||||
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();
|
||||
}
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -342,7 +342,7 @@ void RegisterShowcaseRoutes() {
|
||||
if (!zone) return;
|
||||
const auto asset = PathId<uint32_t>(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<uint32_t>(context.path, 5);
|
||||
const auto slot = PathId<uint32_t>(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));
|
||||
});
|
||||
}
|
||||
|
||||
140
dDashboardServer/routes/WorkerPool.cpp
Normal file
140
dDashboardServer/routes/WorkerPool.cpp
Normal file
@@ -0,0 +1,140 @@
|
||||
#include "WorkerPool.h"
|
||||
|
||||
#include <algorithm>
|
||||
#include <exception>
|
||||
|
||||
#include "Game.h"
|
||||
#include "Logger.h"
|
||||
|
||||
std::optional<WorkerPool::ePriority> WorkerPool::Pick(const std::array<size_t, PRIORITIES>& queued, bool fastLane, size_t runningBackground, size_t maxBackground) {
|
||||
for (size_t i = 0; i < PRIORITIES; i++) {
|
||||
const auto priority = static_cast<ePriority>(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<size_t>(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<size_t>(1, maxBackground);
|
||||
threads = std::max<size_t>(2, threads);
|
||||
for (size_t i = 0; i < threads; i++) m_Threads.emplace_back([this, i] { Work(i == 0); });
|
||||
}
|
||||
|
||||
void WorkerPool::Stop() {
|
||||
std::vector<std::thread> 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<size_t>(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<ePriority> priority;
|
||||
m_Wake.wait(lock, [&] {
|
||||
if (m_Stopping) return true;
|
||||
std::array<size_t, PRIORITIES> 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<size_t>(*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();
|
||||
}
|
||||
}
|
||||
95
dDashboardServer/routes/WorkerPool.h
Normal file
95
dDashboardServer/routes/WorkerPool.h
Normal file
@@ -0,0 +1,95 @@
|
||||
#pragma once
|
||||
|
||||
#include <array>
|
||||
#include <condition_variable>
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <deque>
|
||||
#include <functional>
|
||||
#include <mutex>
|
||||
#include <optional>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
/**
|
||||
* 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<void()>;
|
||||
|
||||
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<ePriority> Pick(const std::array<size_t, PRIORITIES>& 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<std::deque<Entry>, PRIORITIES> m_Queues;
|
||||
std::vector<std::thread> m_Threads;
|
||||
size_t m_MaxBackground{ 1 };
|
||||
size_t m_RunningBackground{};
|
||||
size_t m_Active{};
|
||||
bool m_Stopping{};
|
||||
};
|
||||
@@ -1,5 +1,6 @@
|
||||
set(DWEB_SOURCES
|
||||
"Web.cpp")
|
||||
"Web.cpp"
|
||||
"DeferredReply.cpp")
|
||||
|
||||
add_library(dWeb STATIC ${DWEB_SOURCES})
|
||||
|
||||
|
||||
69
dWeb/DeferredReply.cpp
Normal file
69
dWeb/DeferredReply.cpp
Normal file
@@ -0,0 +1,69 @@
|
||||
#include "DeferredReply.h"
|
||||
|
||||
#include <unordered_map>
|
||||
|
||||
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<DeferredState> DeferredQueue::Begin(unsigned long connection) {
|
||||
auto state = std::make_shared<DeferredState>();
|
||||
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<std::string> 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::Finished> DeferredQueue::Drain() {
|
||||
std::vector<std::pair<uint64_t, HTTPReply>> done;
|
||||
{
|
||||
std::lock_guard lock(m_Mutex);
|
||||
done.swap(m_Done);
|
||||
}
|
||||
std::vector<Finished> finished;
|
||||
if (done.empty()) return finished;
|
||||
std::unordered_map<uint64_t, unsigned long> 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));
|
||||
}
|
||||
93
dWeb/DeferredReply.h
Normal file
93
dWeb/DeferredReply.h
Normal file
@@ -0,0 +1,93 @@
|
||||
#pragma once
|
||||
|
||||
#include <atomic>
|
||||
#include <cstdint>
|
||||
#include <map>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
#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<bool> cancelled{};
|
||||
std::atomic<bool> 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<DeferredState> 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<DeferredState> 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<DeferredState> 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<std::string> 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<Finished> 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<DeferredState> state;
|
||||
std::vector<std::string> headers;
|
||||
bool close{};
|
||||
};
|
||||
|
||||
std::map<unsigned long, Waiting> m_Pending; // web thread only
|
||||
uint64_t m_NextId{ 1 };
|
||||
|
||||
std::mutex m_Mutex;
|
||||
std::vector<std::pair<uint64_t, HTTPReply>> m_Done; // guarded by m_Mutex, keyed by DeferredState::id
|
||||
};
|
||||
@@ -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;
|
||||
|
||||
37
dWeb/HTTPReply.h
Normal file
37
dWeb/HTTPReply.h
Normal file
@@ -0,0 +1,37 @@
|
||||
#pragma once
|
||||
|
||||
#include <memory>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
#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<std::string> 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<DeferredState> deferred{};
|
||||
};
|
||||
133
dWeb/Web.cpp
133
dWeb/Web.cpp
@@ -72,6 +72,9 @@ namespace {
|
||||
|
||||
// Global middleware applied to all routes
|
||||
std::vector<MiddlewarePtr> 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 "<file>.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<int>(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 "<file>.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<mg_http_message*>(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<int>(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 = {};
|
||||
|
||||
46
dWeb/Web.h
46
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<std::string> 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<uint8_t()> 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<bool(uint8_t)> 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<bool(uint8_t)>& 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<bool(uint8_t)> wsApiAccessCallback = nullptr;
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
194
tests/dWebTests/DeferredReplyTests.cpp
Normal file
194
tests/dWebTests/DeferredReplyTests.cpp
Normal file
@@ -0,0 +1,194 @@
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <functional>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
#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<DeferredReply> replies;
|
||||
for (unsigned long i = 1; i <= COUNT; i++) replies.emplace_back(queue.Begin(i));
|
||||
std::vector<std::thread> 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<std::string> bodies;
|
||||
bool closed{};
|
||||
};
|
||||
|
||||
void ClientEvents(mg_connection* connection, int event, void* data) {
|
||||
auto* client = static_cast<Client*>(connection->fn_data);
|
||||
if (event == MG_EV_HTTP_MSG) {
|
||||
const auto* message = static_cast<mg_http_message*>(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<bool> release{};
|
||||
std::vector<std::thread> workers;
|
||||
std::vector<DeferredReply> 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<bool()>& 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);
|
||||
}
|
||||
134
tests/dWebTests/WorkerPoolTests.cpp
Normal file
134
tests/dWebTests/WorkerPoolTests.cpp
Normal file
@@ -0,0 +1,134 @@
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
#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<int> count{};
|
||||
for (int i = 0; i < 200; i++) pool.Submit(static_cast<ePriority>(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<bool> 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<std::string> order;
|
||||
std::atomic<bool> 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<std::string>{ "normal 0", "normal 1", "normal 2", "large", "background" }));
|
||||
}
|
||||
|
||||
TEST(WorkerPoolTests, CancelDropsAGroupsQueuedJobs) {
|
||||
WorkerPool pool;
|
||||
pool.Start(2);
|
||||
std::atomic<bool> release{}, blockerStarted{};
|
||||
std::atomic<int> 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<bool> release{}, blockerStarted{};
|
||||
std::atomic<int> 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());
|
||||
}
|
||||
Reference in New Issue
Block a user