summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorisanae <14251494+isanae@users.noreply.github.com>2020-02-14 01:28:50 -0500
committerisanae <14251494+isanae@users.noreply.github.com>2020-02-18 17:25:03 -0500
commitc8fc7abade6e28507f2a1590007a098e647026a9 (patch)
tree13eaab40225ec292633c005f5e49b76ec5bb3f3d /src
parent763a5d6c08006c319ed92f4088a4d3c211f80cf6 (diff)
thread-safe OriginConnection
ThreadPool now keeps threads running keep ModThreads around so avoid reallocating buffers
Diffstat (limited to 'src')
-rw-r--r--src/directoryrefresher.cpp37
-rw-r--r--src/envfs.cpp24
-rw-r--r--src/envfs.h96
-rw-r--r--src/shared/directoryentry.cpp101
-rw-r--r--src/shared/directoryentry.h8
5 files changed, 202 insertions, 64 deletions
diff --git a/src/directoryrefresher.cpp b/src/directoryrefresher.cpp
index 20d6a52b..a59ddf9d 100644
--- a/src/directoryrefresher.cpp
+++ b/src/directoryrefresher.cpp
@@ -203,8 +203,8 @@ struct ModThread
std::wstring modName;
std::wstring path;
int prio = -1;
- env::Directory* dir = nullptr;
DirectoryStats* stats = nullptr;
+ env::DirectoryWalker walker;
std::condition_variable cv;
std::mutex mutex;
@@ -212,7 +212,11 @@ struct ModThread
void wakeup()
{
- ready = true;
+ {
+ std::scoped_lock lock(mutex);
+ ready = true;
+ }
+
cv.notify_one();
}
@@ -222,7 +226,7 @@ struct ModThread
cv.wait(lock, [&]{ return ready; });
SetThisThreadName(QString::fromStdWString(modName + L" refresher"));
- ds->addFromOrigin(modName, path, prio, *stats);
+ ds->addFromOrigin(walker, modName, path, prio, *stats);
/*if (Settings::instance().archiveParsing()) {
addModBSAToStructure(
@@ -237,6 +241,9 @@ struct ModThread
}
};
+env::ThreadPool<ModThread> g_threads;
+
+
void dumpStats(std::vector<DirectoryStats>& stats)
{
static int run = 0;
@@ -268,35 +275,33 @@ void DirectoryRefresher::addMultipleModsFilesToStructure(
MOShared::DirectoryEntry *directoryStructure,
const std::vector<EntryInfo>& entries, bool emitProgress)
{
- std::vector<env::Directory> dirs(entries.size());
std::vector<DirectoryStats> stats(entries.size());
+ g_threads.setMax(m_threadCount);
+
{
TimeThis tt("walk dirs");
- env::ThreadPool<ModThread> threads(m_threadCount);
-
for (std::size_t i=0; i<entries.size(); ++i) {
const auto& e = entries[i];
const int prio = static_cast<int>(i + 1);
+ stats[i].mod = entries[i].modName.toStdString();
+
try
{
if (e.stealFiles.length() > 0) {
stealModFilesIntoStructure(
directoryStructure, e.modName, prio, e.absolutePath, e.stealFiles);
} else {
- auto& mt = threads.request();
+ auto& mt = g_threads.request();
mt.ds = directoryStructure;
mt.modName = entries[i].modName.toStdWString();
mt.path = QDir::toNativeSeparators(e.absolutePath).toStdWString();
mt.prio = prio;
- mt.dir = &dirs[i];
mt.stats = &stats[i];
- stats[i].mod = entries[i].modName.toStdString();
-
mt.wakeup();
}
} catch (const std::exception& ex) {
@@ -308,7 +313,7 @@ void DirectoryRefresher::addMultipleModsFilesToStructure(
}
}
- threads.join();
+ g_threads.waitForAll();
}
dumpStats(stats);
@@ -347,12 +352,12 @@ void DirectoryRefresher::refresh()
m_DirectoryStructure->getFileRegister()->sortOrigins();
- emit progress(100);
-
cleanStructure(m_DirectoryStructure);
+ }
- emit refreshed();
+ emit progress(100);
- //logcounts("after refresh");
- }
+ emit refreshed();
+
+ //logcounts("after refresh");
}
diff --git a/src/envfs.cpp b/src/envfs.cpp
index 5cd36957..a67faf84 100644
--- a/src/envfs.cpp
+++ b/src/envfs.cpp
@@ -212,15 +212,6 @@ void setHandleCloserThreadCount(std::size_t n)
g_handleClosers.setMax(n);
}
-void shrinkFs()
-{
- g_handleClosers.join();
-
- g_handleClosers.forEach([](auto&& t) {
- t.shrink();
- });
-}
-
void forEachEntryImpl(
void* cx, HandleCloserThread& hc, std::vector<std::unique_ptr<unsigned char[]>>& buffers,
POBJECT_ATTRIBUTES poa, std::size_t depth,
@@ -328,14 +319,13 @@ void forEachEntryImpl(
}
}
-void forEachEntry(
+
+void DirectoryWalker::forEachEntry(
const std::wstring& path, void* cx,
DirStartF* dirStartF, DirEndF* dirEndF, FileF* fileF)
{
auto& hc = g_handleClosers.request();
- std::vector<std::unique_ptr<unsigned char[]>> buffers;
-
if (!NtOpenFile) {
LibraryPtr m(::LoadLibraryW(L"ntdll.dll"));
NtOpenFile = (NtOpenFile_type)::GetProcAddress(m.get(), "NtOpenFile");
@@ -354,10 +344,18 @@ void forEachEntry(
oa.Length = sizeof(oa);
oa.ObjectName = &ObjectName;
- forEachEntryImpl(cx, hc, buffers, &oa, 0, dirStartF, dirEndF, fileF);
+ forEachEntryImpl(cx, hc, m_buffers, &oa, 0, dirStartF, dirEndF, fileF);
hc.wakeup();
}
+
+void forEachEntry(
+ const std::wstring& path, void* cx,
+ DirStartF* dirStartF, DirEndF* dirEndF, FileF* fileF)
+{
+ DirectoryWalker().forEachEntry(path, cx, dirStartF, dirEndF, fileF);
+}
+
Directory getFilesAndDirs(const std::wstring& path)
{
struct Context
diff --git a/src/envfs.h b/src/envfs.h
index 001f2d09..8790b071 100644
--- a/src/envfs.h
+++ b/src/envfs.h
@@ -33,13 +33,13 @@ class ThreadPool
{
public:
ThreadPool(std::size_t max=1)
- : m_threads(max)
{
+ setMax(max);
}
~ThreadPool()
{
- join();
+ stopAndJoin();
}
void setMax(std::size_t n)
@@ -47,15 +47,40 @@ public:
m_threads.resize(n);
}
- void join()
+ void stopAndJoin()
{
for (auto& ti : m_threads) {
+ ti.stop = true;
+ ti.wakeup();
+ }
+
+ for (auto& ti : m_threads) {
if (ti.thread.joinable()) {
ti.thread.join();
}
}
}
+ void waitForAll()
+ {
+ for (;;) {
+ bool done = true;
+
+ for (auto& ti : m_threads) {
+ if (ti.busy) {
+ done = false;
+ break;
+ }
+ }
+
+ if (done) {
+ break;
+ }
+
+ std::this_thread::sleep_for(std::chrono::milliseconds(1));
+ }
+ }
+
T& request()
{
if (m_threads.empty()) {
@@ -67,15 +92,7 @@ public:
bool expected = false;
if (ti.busy.compare_exchange_strong(expected, true)) {
- if (ti.thread.joinable()) {
- ti.thread.join();
- }
-
- ti.thread = std::thread([&]{
- ti.o.run();
- ti.busy = false;
- });
-
+ ti.wakeup();
return ti.o;
}
}
@@ -98,6 +115,47 @@ private:
std::thread thread;
std::atomic<bool> busy;
T o;
+
+ std::condition_variable cv;
+ std::mutex mutex;
+ bool ready;
+
+ std::atomic<bool> stop;
+
+ ThreadInfo()
+ : busy(true), ready(false), stop(false)
+ {
+ thread = std::thread([&]{ run(); });
+ }
+
+ void wakeup()
+ {
+ {
+ std::scoped_lock lock(mutex);
+ ready = true;
+ }
+
+ cv.notify_one();
+ }
+
+ void run()
+ {
+ busy = false;
+
+ while (!stop) {
+ std::unique_lock lock(mutex);
+ cv.wait(lock, [&]{ return ready; });
+
+ if (stop) {
+ break;
+ }
+
+ o.run();
+
+ ready = false;
+ busy = false;
+ }
+ }
};
std::list<ThreadInfo> m_threads;
@@ -109,7 +167,19 @@ using DirEndF = void (void*, std::wstring_view);
using FileF = void (void*, std::wstring_view, FILETIME);
void setHandleCloserThreadCount(std::size_t n);
-void shrinkFs();
+
+
+class DirectoryWalker
+{
+public:
+ void forEachEntry(
+ const std::wstring& path, void* cx,
+ DirStartF* dirStartF, DirEndF* dirEndF, FileF* fileF);
+
+private:
+ std::vector<std::unique_ptr<unsigned char[]>> m_buffers;
+};
+
void forEachEntry(
const std::wstring& path, void* cx,
diff --git a/src/shared/directoryentry.cpp b/src/shared/directoryentry.cpp
index 19500167..4e64d1c1 100644
--- a/src/shared/directoryentry.cpp
+++ b/src/shared/directoryentry.cpp
@@ -211,35 +211,57 @@ public:
--OriginConnectionCount;
}
+ std::pair<FilesOrigin&, bool> getOrCreate(
+ const std::wstring &originName, const std::wstring &directory, int priority,
+ const boost::shared_ptr<FileRegister>& fileRegister,
+ const boost::shared_ptr<OriginConnection>& originConnection,
+ DirectoryStats& stats)
+ {
+ std::unique_lock lock(m_Mutex);
+
+ auto itor = m_OriginsNameMap.find(originName);
+
+ if (itor == m_OriginsNameMap.end()) {
+ FilesOrigin& origin = createOriginNoLock(
+ originName, directory, priority, fileRegister, originConnection);
+
+ return {origin, true};
+ } else {
+ FilesOrigin& origin = m_Origins[itor->second];
+ lock.unlock();
+
+ origin.enable(true, stats);
+ return {origin, false};
+ }
+ }
+
FilesOrigin& createOrigin(
const std::wstring &originName, const std::wstring &directory, int priority,
boost::shared_ptr<FileRegister> fileRegister,
boost::shared_ptr<OriginConnection> originConnection)
{
- int newID = createID();
-
- auto itor = m_Origins.insert({newID, FilesOrigin(
- newID, originName, directory, priority,
- fileRegister, originConnection)}).first;
-
- m_OriginsNameMap.insert({originName, newID});
- m_OriginsPriorityMap.insert({priority, newID});
+ std::scoped_lock lock(m_Mutex);
- return itor->second;
+ return createOriginNoLock(
+ originName, directory, priority, fileRegister, originConnection);
}
bool exists(const std::wstring &name)
{
+ std::scoped_lock lock(m_Mutex);
return m_OriginsNameMap.find(name) != m_OriginsNameMap.end();
}
FilesOrigin &getByID(Index ID)
{
+ std::scoped_lock lock(m_Mutex);
return m_Origins[ID];
}
const FilesOrigin* findByID(Index ID) const
{
+ std::scoped_lock lock(m_Mutex);
+
auto itor = m_Origins.find(ID);
if (itor == m_Origins.end()) {
@@ -251,6 +273,8 @@ public:
FilesOrigin &getByName(const std::wstring &name)
{
+ std::scoped_lock lock(m_Mutex);
+
std::map<std::wstring, int>::iterator iter = m_OriginsNameMap.find(name);
if (iter != m_OriginsNameMap.end()) {
@@ -264,6 +288,8 @@ public:
void changePriorityLookup(int oldPriority, int newPriority)
{
+ std::scoped_lock lock(m_Mutex);
+
auto iter = m_OriginsPriorityMap.find(oldPriority);
if (iter != m_OriginsPriorityMap.end()) {
@@ -275,6 +301,8 @@ public:
void changeNameLookup(const std::wstring &oldName, const std::wstring &newName)
{
+ std::scoped_lock lock(m_Mutex);
+
auto iter = m_OriginsNameMap.find(oldName);
if (iter != m_OriginsNameMap.end()) {
@@ -291,11 +319,29 @@ private:
std::map<Index, FilesOrigin> m_Origins;
std::map<std::wstring, Index> m_OriginsNameMap;
std::map<int, Index> m_OriginsPriorityMap;
+ mutable std::mutex m_Mutex;
Index createID()
{
return m_NextID++;
}
+
+ FilesOrigin& createOriginNoLock(
+ const std::wstring &originName, const std::wstring &directory, int priority,
+ boost::shared_ptr<FileRegister> fileRegister,
+ boost::shared_ptr<OriginConnection> originConnection)
+ {
+ int newID = createID();
+
+ auto itor = m_Origins.insert({newID, FilesOrigin(
+ newID, originName, directory, priority,
+ fileRegister, originConnection)}).first;
+
+ m_OriginsNameMap.insert({originName, newID});
+ m_OriginsPriorityMap.insert({priority, newID});
+
+ return itor->second;
+ }
};
@@ -899,10 +945,18 @@ void DirectoryEntry::addFromOrigin(
const std::wstring &originName, const std::wstring &directory, int priority,
DirectoryStats& stats)
{
+ env::DirectoryWalker walker;
+ addFromOrigin(walker, originName, directory, priority, stats);
+}
+
+void DirectoryEntry::addFromOrigin(
+ env::DirectoryWalker& walker, const std::wstring &originName,
+ const std::wstring &directory, int priority, DirectoryStats& stats)
+{
FilesOrigin &origin = createOrigin(originName, directory, priority, stats);
if (!directory.empty()) {
- addFiles(origin, directory, stats);
+ addFiles(walker, origin, directory, stats);
}
m_Populated = true;
@@ -993,7 +1047,10 @@ void DirectoryEntry::addFromBSA(
void DirectoryEntry::propagateOrigin(int origin)
{
- m_Origins.insert(origin);
+ {
+ std::scoped_lock lock(m_OriginsMutex);
+ m_Origins.insert(origin);
+ }
if (m_Parent != nullptr) {
m_Parent->propagateOrigin(origin);
@@ -1252,16 +1309,17 @@ FilesOrigin &DirectoryEntry::createOrigin(
const std::wstring &originName, const std::wstring &directory, int priority,
DirectoryStats& stats)
{
- if (m_OriginConnection->exists(originName)) {
- ++stats.originExists;
- FilesOrigin &origin = m_OriginConnection->getByName(originName);
- origin.enable(true, stats);
- return origin;
- } else {
+ auto r = m_OriginConnection->getOrCreate(
+ originName, directory, priority,
+ m_FileRegister, m_OriginConnection, stats);
+
+ if (r.second) {
++stats.originCreate;
- return m_OriginConnection->createOrigin(
- originName, directory, priority, m_FileRegister, m_OriginConnection);
+ } else {
+ ++stats.originExists;
}
+
+ return r.first;
}
void DirectoryEntry::removeFiles(const std::set<FileEntry::Index> &indices)
@@ -1359,7 +1417,8 @@ FileEntry::Ptr DirectoryEntry::insert(
}
void DirectoryEntry::addFiles(
- FilesOrigin &origin, const std::wstring& path, DirectoryStats& stats)
+ env::DirectoryWalker& walker, FilesOrigin &origin,
+ const std::wstring& path, DirectoryStats& stats)
{
struct Context
{
@@ -1371,7 +1430,7 @@ void DirectoryEntry::addFiles(
Context cx = {origin, stats};
cx.current.push(this);
- env::forEachEntry(path, &cx,
+ walker.forEachEntry(path, &cx,
[](void* pcx, std::wstring_view path)
{
Context* cx = (Context*)pcx;
diff --git a/src/shared/directoryentry.h b/src/shared/directoryentry.h
index 8f6afbb1..2a70d486 100644
--- a/src/shared/directoryentry.h
+++ b/src/shared/directoryentry.h
@@ -410,6 +410,10 @@ public:
const std::wstring &originName,
const std::wstring &directory, int priority, DirectoryStats& stats);
+ void addFromOrigin(
+ env::DirectoryWalker& walker, const std::wstring &originName,
+ const std::wstring &directory, int priority, DirectoryStats& stats);
+
void addFromBSA(
const std::wstring &originName, std::wstring &directory,
const std::wstring &fileName, int priority, int order);
@@ -556,6 +560,7 @@ private:
bool m_TopLevel;
std::mutex m_SubDirMutex;
std::mutex m_FilesMutex;
+ std::mutex m_OriginsMutex;
DirectoryEntry(const DirectoryEntry &reference);
@@ -569,7 +574,8 @@ private:
std::wstring_view archive, int order, DirectoryStats& stats);
void addFiles(
- FilesOrigin &origin, const std::wstring& path, DirectoryStats& stats);
+ env::DirectoryWalker& walker, FilesOrigin &origin,
+ const std::wstring& path, DirectoryStats& stats);
void addFiles(
FilesOrigin &origin, BSA::Folder::Ptr archiveFolder, FILETIME &fileTime,