From c8fc7abade6e28507f2a1590007a098e647026a9 Mon Sep 17 00:00:00 2001 From: isanae <14251494+isanae@users.noreply.github.com> Date: Fri, 14 Feb 2020 01:28:50 -0500 Subject: thread-safe OriginConnection ThreadPool now keeps threads running keep ModThreads around so avoid reallocating buffers --- src/directoryrefresher.cpp | 37 +++++++++------- src/envfs.cpp | 24 +++++----- src/envfs.h | 96 +++++++++++++++++++++++++++++++++------ src/shared/directoryentry.cpp | 101 +++++++++++++++++++++++++++++++++--------- src/shared/directoryentry.h | 8 +++- 5 files changed, 202 insertions(+), 64 deletions(-) (limited to 'src') 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 g_threads; + + void dumpStats(std::vector& stats) { static int run = 0; @@ -268,35 +275,33 @@ void DirectoryRefresher::addMultipleModsFilesToStructure( MOShared::DirectoryEntry *directoryStructure, const std::vector& entries, bool emitProgress) { - std::vector dirs(entries.size()); std::vector stats(entries.size()); + g_threads.setMax(m_threadCount); + { TimeThis tt("walk dirs"); - env::ThreadPool threads(m_threadCount); - for (std::size_t i=0; i(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>& 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> 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,8 +47,13 @@ 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(); @@ -56,6 +61,26 @@ public: } } + 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 busy; T o; + + std::condition_variable cv; + std::mutex mutex; + bool ready; + + std::atomic 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 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> 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 getOrCreate( + const std::wstring &originName, const std::wstring &directory, int priority, + const boost::shared_ptr& fileRegister, + const boost::shared_ptr& 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, boost::shared_ptr 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::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 m_Origins; std::map m_OriginsNameMap; std::map 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, + boost::shared_ptr 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; + } }; @@ -898,11 +944,19 @@ void DirectoryEntry::clear() 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 &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, -- cgit v1.3.1