diff options
| author | isanae <14251494+isanae@users.noreply.github.com> | 2020-02-14 01:28:50 -0500 |
|---|---|---|
| committer | isanae <14251494+isanae@users.noreply.github.com> | 2020-02-18 17:25:03 -0500 |
| commit | c8fc7abade6e28507f2a1590007a098e647026a9 (patch) | |
| tree | 13eaab40225ec292633c005f5e49b76ec5bb3f3d | |
| parent | 763a5d6c08006c319ed92f4088a4d3c211f80cf6 (diff) | |
thread-safe OriginConnection
ThreadPool now keeps threads running
keep ModThreads around so avoid reallocating buffers
| -rw-r--r-- | src/directoryrefresher.cpp | 37 | ||||
| -rw-r--r-- | src/envfs.cpp | 24 | ||||
| -rw-r--r-- | src/envfs.h | 96 | ||||
| -rw-r--r-- | src/shared/directoryentry.cpp | 101 | ||||
| -rw-r--r-- | src/shared/directoryentry.h | 8 |
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,
|
