From ed1be02aafabb63dd15a63840e94c1dcf616fea5 Mon Sep 17 00:00:00 2001 From: emeric Date: Sun, 17 Mar 2024 23:36:40 +0100 Subject: [PATCH] Added a scan tracker to speed up queries with high offset in very large database --- src/libs/core/impl/ChildProcessManager.cpp | 6 +- src/libs/database/impl/Track.cpp | 19 +-- src/libs/database/include/database/Track.hpp | 2 +- src/libs/database/test/Track.cpp | 46 +++++- .../impl/ScanStepRemoveOrphanDbFiles.cpp | 20 ++- src/libs/subsonic/impl/ClientInfo.hpp | 1 + src/libs/subsonic/impl/SubsonicResource.cpp | 7 +- src/libs/subsonic/impl/SubsonicResource.hpp | 2 +- .../subsonic/impl/entrypoints/Searching.cpp | 156 +++++++++++++++--- 9 files changed, 196 insertions(+), 63 deletions(-) diff --git a/src/libs/core/impl/ChildProcessManager.cpp b/src/libs/core/impl/ChildProcessManager.cpp index abe6de1e..e7dc459e 100644 --- a/src/libs/core/impl/ChildProcessManager.cpp +++ b/src/libs/core/impl/ChildProcessManager.cpp @@ -25,8 +25,7 @@ namespace lms::core { - std::unique_ptr - createChildProcessManager(boost::asio::io_context& ioContext) + std::unique_ptr createChildProcessManager(boost::asio::io_context& ioContext) { return std::make_unique(ioContext); } @@ -36,8 +35,7 @@ namespace lms::core { } - std::unique_ptr - ChildProcessManager::spawnChildProcess(const std::filesystem::path& path, const IChildProcess::Args& args) + std::unique_ptr ChildProcessManager::spawnChildProcess(const std::filesystem::path& path, const IChildProcess::Args& args) { return std::make_unique(_ioContext, path, args); } diff --git a/src/libs/database/impl/Track.cpp b/src/libs/database/impl/Track.cpp index 16d79960..4ff6ed61 100644 --- a/src/libs/database/impl/Track.cpp +++ b/src/libs/database/impl/Track.cpp @@ -228,27 +228,22 @@ namespace lms::db .resultValue(); } - void Track::find(Session& session, TrackId& lastRetrievedTrack, std::size_t batchSize, bool& moreResults, const std::function& func) + void Track::find(Session& session, TrackId& lastRetrievedTrack, std::size_t count, const std::function& func, MediaLibraryId library) { session.checkReadTransaction(); - auto collection{ session.getDboSession().find() + auto query{ session.getDboSession().find() .orderBy("id") .where("id > ?").bind(lastRetrievedTrack) - .limit(static_cast(batchSize) + 1) - .resultList() }; + .limit(static_cast(count)) }; - moreResults = false; + if (library.isValid()) + query.where("media_library_id = ?").bind(library); + + auto collection{query.resultList()}; - std::size_t count{}; for (auto itResult{ collection.begin() }; itResult != collection.end(); ++itResult) { - if (count++ == batchSize) - { - moreResults = true; - break; - } - func(*itResult); lastRetrievedTrack = (*itResult)->getId(); } diff --git a/src/libs/database/include/database/Track.hpp b/src/libs/database/include/database/Track.hpp index c52f9f49..6b29414a 100644 --- a/src/libs/database/include/database/Track.hpp +++ b/src/libs/database/include/database/Track.hpp @@ -110,7 +110,7 @@ namespace lms::db static std::size_t getCount(Session& session); static pointer findByPath(Session& session, const std::filesystem::path& p); static pointer find(Session& session, TrackId id); - static void find(Session& session, TrackId& lastRetrievedTrack, std::size_t batchSize, bool& moreResults, const std::function& func); + static void find(Session& session, TrackId& lastRetrievedTrack, std::size_t count, const std::function& func, MediaLibraryId library = {}); static bool exists(Session& session, TrackId id); static std::vector findByRecordingMBID(Session& session, const core::UUID& MBID); static std::vector findByMBID(Session& session, const core::UUID& MBID); diff --git a/src/libs/database/test/Track.cpp b/src/libs/database/test/Track.cpp index 0d6683af..c3a0df8f 100644 --- a/src/libs/database/test/Track.cpp +++ b/src/libs/database/test/Track.cpp @@ -68,14 +68,20 @@ namespace lms::db::tests ScopedTrack track1{ session, "MyTrackFile1" }; ScopedTrack track2{ session, "MyTrackFile1" }; ScopedTrack track3{ session, "MyTrackFile1" }; + ScopedMediaLibrary library{ session }; + ScopedMediaLibrary otherLibrary{ session }; + + { + auto transaction{ session.createWriteTransaction() }; + track2.get().modify()->setMediaLibrary(library.get()); + } { auto transaction{ session.createReadTransaction() }; - bool moreResults; TrackId lastRetrievedTrackId; std::vector visitedTracks; - Track::find(session, lastRetrievedTrackId, 10, moreResults, [&](const Track::pointer& track) + Track::find(session, lastRetrievedTrackId, 10, [&](const Track::pointer& track) { visitedTracks.push_back(track); }); @@ -83,40 +89,62 @@ namespace lms::db::tests EXPECT_EQ(visitedTracks[0]->getId(), track1.getId()); EXPECT_EQ(visitedTracks[1]->getId(), track2.getId()); EXPECT_EQ(visitedTracks[2]->getId(), track3.getId()); - EXPECT_FALSE(moreResults); EXPECT_EQ(lastRetrievedTrackId, track3.getId()); } { auto transaction{ session.createReadTransaction() }; - bool moreResults; TrackId lastRetrievedTrackId{ track1.getId() }; std::vector visitedTracks; - Track::find(session, lastRetrievedTrackId, 1, moreResults, [&](const Track::pointer& track) + Track::find(session, lastRetrievedTrackId, 1, [&](const Track::pointer& track) { visitedTracks.push_back(track); }); ASSERT_EQ(visitedTracks.size(), 1); EXPECT_EQ(visitedTracks[0]->getId(), track2.getId()); - EXPECT_TRUE(moreResults); EXPECT_EQ(lastRetrievedTrackId, track2.getId()); } { auto transaction{ session.createReadTransaction() }; - bool moreResults; TrackId lastRetrievedTrackId{ track1.getId() }; std::vector visitedTracks; - Track::find(session, lastRetrievedTrackId, 0, moreResults, [&](const Track::pointer& track) + Track::find(session, lastRetrievedTrackId, 0, [&](const Track::pointer& track) { visitedTracks.push_back(track); }); ASSERT_EQ(visitedTracks.size(), 0); - EXPECT_TRUE(moreResults); EXPECT_EQ(lastRetrievedTrackId, track1.getId()); } + + { + auto transaction{ session.createReadTransaction() }; + + TrackId lastRetrievedTrackId{}; + std::vector visitedTracks; + Track::find(session, lastRetrievedTrackId, 10, [&](const Track::pointer& track) + { + visitedTracks.push_back(track); + }, otherLibrary.getId()); + ASSERT_EQ(visitedTracks.size(), 0); + EXPECT_EQ(lastRetrievedTrackId, TrackId{}); + } + + { + auto transaction{ session.createReadTransaction() }; + + TrackId lastRetrievedTrackId{}; + std::vector visitedTracks; + Track::find(session, lastRetrievedTrackId, 10, [&](const Track::pointer& track) + { + visitedTracks.push_back(track); + }, library.getId()); + ASSERT_EQ(visitedTracks.size(), 1); + EXPECT_EQ(visitedTracks[0]->getId(), track2.getId()); + EXPECT_EQ(lastRetrievedTrackId, track2.getId()); + } } TEST_F(DatabaseFixture, Track_MediaLibrary) diff --git a/src/libs/services/scanner/impl/ScanStepRemoveOrphanDbFiles.cpp b/src/libs/services/scanner/impl/ScanStepRemoveOrphanDbFiles.cpp index 00ad637a..a0a30185 100644 --- a/src/libs/services/scanner/impl/ScanStepRemoveOrphanDbFiles.cpp +++ b/src/libs/services/scanner/impl/ScanStepRemoveOrphanDbFiles.cpp @@ -99,8 +99,8 @@ namespace lms::scanner std::vector tracksToRemove; TrackId lastCheckedTrackID; - bool moreResults{ true }; - while (moreResults) + bool endReached{}; + while (!endReached) { if (_abortScan) break; @@ -108,13 +108,17 @@ namespace lms::scanner tracksToRemove.clear(); { auto transaction{ session.createReadTransaction() }; - Track::find(session, lastCheckedTrackID, batchSize, moreResults, [&](const Track::pointer& track) - { - if (!checkFile(track->getPath())) - tracksToRemove.push_back(track); - context.currentStepStats.processedElems++; - }); + endReached = true; + Track::find(session, lastCheckedTrackID, batchSize, [&](const Track::pointer& track) + { + endReached = false; + + if (!checkFile(track->getPath())) + tracksToRemove.push_back(track); + + context.currentStepStats.processedElems++; + }); } if (!tracksToRemove.empty()) diff --git a/src/libs/subsonic/impl/ClientInfo.hpp b/src/libs/subsonic/impl/ClientInfo.hpp index beec15ba..045bafcf 100644 --- a/src/libs/subsonic/impl/ClientInfo.hpp +++ b/src/libs/subsonic/impl/ClientInfo.hpp @@ -26,6 +26,7 @@ namespace lms::api::subsonic { struct ClientInfo { + std::string ipAddress; std::string name; std::string user; std::string password; diff --git a/src/libs/subsonic/impl/SubsonicResource.cpp b/src/libs/subsonic/impl/SubsonicResource.cpp index e1c232b8..e709a96b 100644 --- a/src/libs/subsonic/impl/SubsonicResource.cpp +++ b/src/libs/subsonic/impl/SubsonicResource.cpp @@ -380,13 +380,16 @@ namespace lms::api::subsonic } } - ClientInfo SubsonicResource::getClientInfo(const Wt::Http::ParameterMap& parameters) + ClientInfo SubsonicResource::getClientInfo(const Wt::Http::Request& request) { + const auto& parameters{ request.getParameterMap() }; ClientInfo res; if (hasParameter(parameters, "t")) throw TokenAuthenticationNotSupportedForLDAPUsersError{}; + res.ipAddress = request.clientAddress(); + // Mandatory parameters res.name = getMandatoryParameterAs(parameters, "c"); res.version = getMandatoryParameterAs(parameters, "v"); @@ -399,7 +402,7 @@ namespace lms::api::subsonic RequestContext SubsonicResource::buildRequestContext(const Wt::Http::Request& request) { const Wt::Http::ParameterMap& parameters{ request.getParameterMap() }; - const ClientInfo clientInfo{ getClientInfo(parameters) }; + const ClientInfo clientInfo{ getClientInfo(request) }; const db::UserId userId{ authenticateUser(request, clientInfo) }; bool enableOpenSubsonic{ _openSubsonicDisabledClients.find(clientInfo.name) == std::cend(_openSubsonicDisabledClients) }; bool enableDefaultCover{ _defaultCoverClients.find(clientInfo.name) != std::cend(_openSubsonicDisabledClients) }; diff --git a/src/libs/subsonic/impl/SubsonicResource.hpp b/src/libs/subsonic/impl/SubsonicResource.hpp index eda68edf..67db9936 100644 --- a/src/libs/subsonic/impl/SubsonicResource.hpp +++ b/src/libs/subsonic/impl/SubsonicResource.hpp @@ -47,7 +47,7 @@ namespace lms::api::subsonic ProtocolVersion getServerProtocolVersion(const std::string& clientName) const; static void checkProtocolVersion(ProtocolVersion client, ProtocolVersion server); - ClientInfo getClientInfo(const Wt::Http::ParameterMap& parameters); + ClientInfo getClientInfo(const Wt::Http::Request& request); RequestContext buildRequestContext(const Wt::Http::Request& request); db::UserId authenticateUser(const Wt::Http::Request& request, const ClientInfo& clientInfo); diff --git a/src/libs/subsonic/impl/entrypoints/Searching.cpp b/src/libs/subsonic/impl/entrypoints/Searching.cpp index ae876f6b..a8669b99 100644 --- a/src/libs/subsonic/impl/entrypoints/Searching.cpp +++ b/src/libs/subsonic/impl/entrypoints/Searching.cpp @@ -19,6 +19,10 @@ #include "Searching.hpp" +#include +#include +#include + #include "database/Artist.hpp" #include "database/Release.hpp" #include "database/Session.hpp" @@ -29,42 +33,111 @@ #include "responses/Song.hpp" #include "ParameterParsing.hpp" #include "SubsonicId.hpp" - + namespace lms::api::subsonic { using namespace db; + namespace + { + // Search endpoints can be used to scan/sync the database + // This class is used to keep track of the current scans, in order to retrieve the last objectId + // to speed up the query of the following range (avoid the 'offset' cost) + template + class ScanTracker + { + public: + ObjectId extractLastRetrievedObjectId(const ClientInfo& info, std::size_t offset); + void setObjectId(const ClientInfo& info, std::size_t offset, ObjectId lastRetrievedId); + void cleanOutdatedScanEntries(); + + private: + using ClockType = std::chrono::steady_clock; + + struct Client + { + std::string clientAddress; + std::string clientName; + std::string userName; + std::size_t offset{}; + auto operator<=>(const Client&) const = default; + }; + struct Entry + { + ClockType::time_point timePoint; + ObjectId objectId; + }; + + static constexpr ClockType::duration maxEntryDuration{ std::chrono::seconds{30} }; + + std::mutex _mutex; + std::map _ongoingScans; + }; + + template + ObjectId ScanTracker::extractLastRetrievedObjectId(const ClientInfo& info, std::size_t offset) + { + ObjectId res; + + const std::scoped_lock lock{ _mutex }; + auto it{ _ongoingScans.find({ info.ipAddress, info.name, info.user, offset }) }; + if (it != _ongoingScans.end()) + { + res = it->second.objectId; + _ongoingScans.erase(it); + } + + return res; + } + + template + void ScanTracker::setObjectId(const ClientInfo& info, std::size_t offset, ObjectId lastRetrievedId) + { + const std::scoped_lock lock{ _mutex }; + _ongoingScans[Client{ info.ipAddress, info.name, info.user, offset }] = { ClockType::now(), lastRetrievedId }; + } + + template + void ScanTracker::cleanOutdatedScanEntries() + { + const ClockType::time_point now{ ClockType::now() }; + + const std::scoped_lock lock{ _mutex }; + std::erase_if(_ongoingScans, [&](const auto& entry) { return now > entry.second.timePoint + maxEntryDuration; }); + } + } + namespace { Response handleSearchRequestCommon(RequestContext& context, bool id3) { // Mandatory params - std::string queryString{ getMandatoryParameterAs(context.parameters, "query") }; + const std::string queryString{ getMandatoryParameterAs(context.parameters, "query") }; std::string_view query{ queryString }; // Optional params + const std::size_t artistCount{ getParameterAs(context.parameters, "artistCount").value_or(20) }; + const std::size_t artistOffset{ getParameterAs(context.parameters, "artistOffset").value_or(0) }; + const std::size_t albumCount{ getParameterAs(context.parameters, "albumCount").value_or(20) }; + const std::size_t albumOffset{ getParameterAs(context.parameters, "albumOffset").value_or(0) }; + const std::size_t songCount{ getParameterAs(context.parameters, "songCount").value_or(20) }; + const std::size_t songOffset{ getParameterAs(context.parameters, "songOffset").value_or(0) }; const MediaLibraryId mediaLibrary{ getParameterAs(context.parameters, "musicFolderId").value_or(MediaLibraryId{}) }; + if (artistCount > defaultMaxCountSize) + throw ParameterValueTooHighGenericError{ "artistCount", defaultMaxCountSize }; + if (albumCount > defaultMaxCountSize) + throw ParameterValueTooHighGenericError{ "albumCount", defaultMaxCountSize }; + if (songCount > defaultMaxCountSize) + throw ParameterValueTooHighGenericError{ "songCount", defaultMaxCountSize }; + // Symfonium adds extra "" if (context.clientInfo.name == "Symfonium") query = core::stringUtils::stringTrim(query, "\""); - std::vector keywords{ core::stringUtils::splitString(query, ' ') }; - - // Optional params - std::size_t artistCount{ getParameterAs(context.parameters, "artistCount").value_or(20) }; - std::size_t artistOffset{ getParameterAs(context.parameters, "artistOffset").value_or(0) }; - std::size_t albumCount{ getParameterAs(context.parameters, "albumCount").value_or(20) }; - std::size_t albumOffset{ getParameterAs(context.parameters, "albumOffset").value_or(0) }; - std::size_t songCount{ getParameterAs(context.parameters, "songCount").value_or(20) }; - std::size_t songOffset{ getParameterAs(context.parameters, "songOffset").value_or(0) }; - - if (artistCount > defaultMaxCountSize) - throw ParameterValueTooHighGenericError{ "artistCount", defaultMaxCountSize }; - else if (albumCount > defaultMaxCountSize) - throw ParameterValueTooHighGenericError{ "albumCount", defaultMaxCountSize }; - else if (songCount > defaultMaxCountSize) - throw ParameterValueTooHighGenericError{ "songCount", defaultMaxCountSize }; + std::vector keywords; + if (!query.empty()) + keywords = core::stringUtils::splitString(query, ' '); Response response{ Response::createOkResponse(context.serverProtocolVersion) }; Response::Node& searchResult2Node{ response.createNode(id3 ? "searchResult3" : "searchResult2") }; @@ -103,15 +176,47 @@ namespace lms::api::subsonic if (songCount > 0) { - Track::FindParameters params; - params.setKeywords(keywords); - params.setRange(Range{ songOffset, songCount }); - params.setMediaLibrary(mediaLibrary); + static ScanTracker currentScansInProgress; + currentScansInProgress.cleanOutdatedScanEntries(); - Track::find(context.dbSession, params, [&](const Track::pointer& track) + TrackId lastRetrievedId; + + auto findTracks{ [&] + { + Track::FindParameters params; + params.setKeywords(keywords); + params.setRange(Range{ songOffset, songCount }); + params.setMediaLibrary(mediaLibrary); + + Track::find(context.dbSession, params, [&](const Track::pointer& track) + { + searchResult2Node.addArrayChild("song", createSongNode(context, track, user)); + lastRetrievedId = track->getId(); + }); + } }; + + if (!keywords.empty()) + { + findTracks(); + } + else + { + if (TrackId cachedLastRetrievedId{ currentScansInProgress.extractLastRetrievedObjectId(context.clientInfo, songOffset) }; cachedLastRetrievedId.isValid()) { - searchResult2Node.addArrayChild("song", createSongNode(context, track, user)); - }); + Track::find(context.dbSession, cachedLastRetrievedId, songCount, [&](const Track::pointer& track) + { + searchResult2Node.addArrayChild("song", createSongNode(context, track, user)); + }, mediaLibrary); + lastRetrievedId = cachedLastRetrievedId; + } + else + { + findTracks(); + } + + if (lastRetrievedId.isValid()) + currentScansInProgress.setObjectId(context.clientInfo, songOffset + songCount, lastRetrievedId); + } } return response; @@ -127,5 +232,4 @@ namespace lms::api::subsonic { return handleSearchRequestCommon(context, true /* id3 */); } - } \ No newline at end of file