Loading of the recommendation engine now controlled by scanner. Bonus: better control/reporting
This commit is contained in:
@@ -19,16 +19,39 @@
|
||||
|
||||
#include "Engine.hpp"
|
||||
|
||||
#include "recommendation/ClustersClassifierCreator.hpp"
|
||||
#include "recommendation/FeaturesClassifierCreator.hpp"
|
||||
#include <unordered_map>
|
||||
#include <vector>
|
||||
|
||||
#include "ClustersClassifierCreator.hpp"
|
||||
#include "FeaturesClassifierCreator.hpp"
|
||||
|
||||
#include "database/Db.hpp"
|
||||
#include "database/Session.hpp"
|
||||
#include "database/ScanSettings.hpp"
|
||||
#include "database/TrackList.hpp"
|
||||
#include "utils/Exception.hpp"
|
||||
#include "utils/Logger.hpp"
|
||||
|
||||
namespace Recommendation {
|
||||
|
||||
|
||||
static
|
||||
std::unique_ptr<IClassifier>
|
||||
createClassifier(ClassifierType type)
|
||||
{
|
||||
switch (type)
|
||||
{
|
||||
case ClassifierType::Clusters:
|
||||
return createClustersClassifier();
|
||||
break;
|
||||
|
||||
case ClassifierType::Features:
|
||||
return createFeaturesClassifier();
|
||||
break;
|
||||
}
|
||||
|
||||
return {};
|
||||
}
|
||||
|
||||
std::unique_ptr<IEngine>
|
||||
createEngine(Database::Db& db)
|
||||
{
|
||||
@@ -36,65 +59,16 @@ createEngine(Database::Db& db)
|
||||
}
|
||||
|
||||
Engine::Engine(Database::Db& db)
|
||||
: _dbSession {db}
|
||||
: _db {db}
|
||||
{
|
||||
start();
|
||||
}
|
||||
|
||||
Engine::~Engine()
|
||||
{
|
||||
stop();
|
||||
}
|
||||
|
||||
void
|
||||
Engine::start()
|
||||
{
|
||||
assert(!_running);
|
||||
_running = true;
|
||||
_ioService.start();
|
||||
}
|
||||
|
||||
void
|
||||
Engine::stop()
|
||||
{
|
||||
assert(_running);
|
||||
_running = false;
|
||||
|
||||
cancelPendingClassifiers();
|
||||
|
||||
_ioService.stop();
|
||||
}
|
||||
|
||||
void
|
||||
Engine::requestLoad()
|
||||
{
|
||||
requestReloadInternal(false);
|
||||
}
|
||||
|
||||
void
|
||||
Engine::requestReload()
|
||||
{
|
||||
requestReloadInternal(true);
|
||||
}
|
||||
|
||||
void
|
||||
Engine::requestReloadInternal(bool databaseChanged)
|
||||
{
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Reload requested...";
|
||||
|
||||
_ioService.post([=]()
|
||||
{
|
||||
reload(databaseChanged);
|
||||
});
|
||||
}
|
||||
|
||||
std::vector<Database::IdType>
|
||||
std::unordered_set<Database::IdType>
|
||||
Engine::getSimilarTracksFromTrackList(Database::Session& session, Database::IdType trackListId, std::size_t maxCount)
|
||||
{
|
||||
std::unordered_set<Database::IdType> res;
|
||||
|
||||
std::shared_lock lock {_classifiersMutex};
|
||||
|
||||
std::vector<Database::IdType> res;
|
||||
|
||||
for (const auto& classifierName : _classifierPriorities)
|
||||
{
|
||||
auto itClassifier {_classifiers.find(classifierName)};
|
||||
@@ -109,23 +83,23 @@ Engine::getSimilarTracksFromTrackList(Database::Session& session, Database::IdTy
|
||||
return res;
|
||||
}
|
||||
|
||||
std::vector<Database::IdType>
|
||||
std::unordered_set<Database::IdType>
|
||||
Engine::getSimilarTracks(Database::Session& dbSession, const std::unordered_set<Database::IdType>& trackIds, std::size_t maxCount)
|
||||
{
|
||||
std::unordered_set<Database::IdType> res;
|
||||
|
||||
std::shared_lock lock {_classifiersMutex};
|
||||
|
||||
std::vector<Database::IdType> res;
|
||||
|
||||
for (const auto& classifierName : _classifierPriorities)
|
||||
for (ClassifierType classifierType : _classifierPriorities)
|
||||
{
|
||||
auto itClassifier {_classifiers.find(classifierName)};
|
||||
auto itClassifier {_classifiers.find(classifierType)};
|
||||
if (itClassifier == std::cend(_classifiers))
|
||||
continue;
|
||||
|
||||
res = itClassifier->second->getSimilarTracks(dbSession, trackIds, maxCount);
|
||||
const IClassifier& classifier {*itClassifier->second};
|
||||
res = classifier.getSimilarTracks(dbSession, trackIds, maxCount);
|
||||
if (!res.empty())
|
||||
{
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar tracks using classifier '" << classifierName << "'";
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar tracks using classifier '" << classifier.getName() << "'";
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -133,23 +107,23 @@ Engine::getSimilarTracks(Database::Session& dbSession, const std::unordered_set<
|
||||
return res;
|
||||
}
|
||||
|
||||
std::vector<Database::IdType>
|
||||
std::unordered_set<Database::IdType>
|
||||
Engine::getSimilarReleases(Database::Session& dbSession, Database::IdType releaseId, std::size_t maxCount)
|
||||
{
|
||||
std::unordered_set<Database::IdType> res;
|
||||
|
||||
std::shared_lock lock {_classifiersMutex};
|
||||
|
||||
std::vector<Database::IdType> res;
|
||||
|
||||
for (const auto& classifierName : _classifierPriorities)
|
||||
for (ClassifierType classifierType : _classifierPriorities)
|
||||
{
|
||||
auto itClassifier {_classifiers.find(classifierName)};
|
||||
auto itClassifier {_classifiers.find(classifierType)};
|
||||
if (itClassifier == std::cend(_classifiers))
|
||||
continue;
|
||||
|
||||
res = itClassifier->second->getSimilarReleases(dbSession, releaseId, maxCount);
|
||||
const IClassifier& classifier {*itClassifier->second};
|
||||
res = classifier.getSimilarReleases(dbSession, releaseId, maxCount);
|
||||
if (!res.empty())
|
||||
{
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar releases using classifier '" << classifierName << "'";
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar releases using classifier '" << classifier.getName() << "'";
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -157,23 +131,23 @@ Engine::getSimilarReleases(Database::Session& dbSession, Database::IdType releas
|
||||
return res;
|
||||
}
|
||||
|
||||
std::vector<Database::IdType>
|
||||
std::unordered_set<Database::IdType>
|
||||
Engine::getSimilarArtists(Database::Session& dbSession, Database::IdType artistId, std::size_t maxCount)
|
||||
{
|
||||
std::unordered_set<Database::IdType> res;
|
||||
|
||||
std::shared_lock lock {_classifiersMutex};
|
||||
|
||||
std::vector<Database::IdType> res;
|
||||
|
||||
for (const auto& classifierName : _classifierPriorities)
|
||||
for (ClassifierType classifierType : _classifierPriorities)
|
||||
{
|
||||
auto itClassifier {_classifiers.find(classifierName)};
|
||||
auto itClassifier {_classifiers.find(classifierType)};
|
||||
if (itClassifier == std::cend(_classifiers))
|
||||
continue;
|
||||
|
||||
res = itClassifier->second->getSimilarArtists(dbSession, artistId, maxCount);
|
||||
const IClassifier& classifier {*itClassifier->second};
|
||||
res = classifier.getSimilarArtists(dbSession, artistId, maxCount);
|
||||
if (!res.empty())
|
||||
{
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar artists using classifier '" << classifierName << "'";
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar artists using classifier '" << classifier.getName() << "'";
|
||||
return res;
|
||||
}
|
||||
}
|
||||
@@ -181,107 +155,132 @@ Engine::getSimilarArtists(Database::Session& dbSession, Database::IdType artistI
|
||||
return res;
|
||||
}
|
||||
|
||||
void
|
||||
Engine::reload(bool databaseChanged)
|
||||
static
|
||||
Database::ScanSettings::RecommendationEngineType
|
||||
getRecommendationEngineType(Database::Session& session)
|
||||
{
|
||||
using namespace Database;
|
||||
auto transaction {session.createSharedTransaction()};
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Reloading recommendation engines...";
|
||||
|
||||
const ScanSettings::RecommendationEngineType engineType {[&]()
|
||||
{
|
||||
auto transaction {_dbSession.createSharedTransaction()};
|
||||
|
||||
return ScanSettings::get(_dbSession)->getRecommendationEngineType();
|
||||
}()};
|
||||
|
||||
clearClassifiers();
|
||||
|
||||
switch (engineType)
|
||||
{
|
||||
case ScanSettings::RecommendationEngineType::Features:
|
||||
{
|
||||
auto clustersClassifier {createClustersClassifier()};
|
||||
auto featuresClassifier {createFeaturesClassifier()};
|
||||
|
||||
setClassifierPriorities({featuresClassifier->getName(), clustersClassifier->getName()});
|
||||
|
||||
initAndAddClassifier(std::move(clustersClassifier), databaseChanged); // init first since faster
|
||||
initAndAddClassifier(std::move(featuresClassifier), databaseChanged);
|
||||
break;
|
||||
}
|
||||
|
||||
case ScanSettings::RecommendationEngineType::Clusters:
|
||||
auto clustersClassifier {createClustersClassifier()};
|
||||
|
||||
setClassifierPriorities({clustersClassifier->getName()});
|
||||
|
||||
initAndAddClassifier(std::move(clustersClassifier), databaseChanged);
|
||||
break;
|
||||
}
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Recommendation engines reloaded!";
|
||||
|
||||
_sigReloaded.emit();
|
||||
return Database::ScanSettings::get(session)->getRecommendationEngineType();
|
||||
}
|
||||
|
||||
void
|
||||
Engine::setClassifierPriorities(std::initializer_list<std::string_view> classifierPriorities)
|
||||
Engine::load(bool forceReload, const ProgressCallback& progressCallback)
|
||||
{
|
||||
using namespace Database;
|
||||
|
||||
static const std::unordered_map<ScanSettings::RecommendationEngineType, std::vector<ClassifierType>> classifierMappings
|
||||
{
|
||||
{ScanSettings::RecommendationEngineType::Features, {ClassifierType::Clusters, ClassifierType::Features}},
|
||||
{ScanSettings::RecommendationEngineType::Clusters, {ClassifierType::Clusters}},
|
||||
};
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Reloading recommendation engines...";
|
||||
|
||||
const ScanSettings::RecommendationEngineType engineType {getRecommendationEngineType(_db.getTLSSession())};
|
||||
|
||||
assert(_pendingClassifiers.empty());
|
||||
clearClassifiers();
|
||||
|
||||
auto itClassifierTypes {classifierMappings.find(engineType)};
|
||||
assert(itClassifierTypes != std::cend(classifierMappings));
|
||||
const std::vector<ClassifierType>& classifierTypes {itClassifierTypes->second};
|
||||
|
||||
setClassifierPriorities(classifierTypes);
|
||||
|
||||
std::vector<std::unique_ptr<IClassifier>> classifiers;
|
||||
for (ClassifierType type : classifierTypes)
|
||||
classifiers.emplace_back(createClassifier(type));
|
||||
|
||||
{
|
||||
std::scoped_lock lock {_controlMutex};
|
||||
|
||||
std::transform(std::cbegin(classifiers), std::cend(classifiers), std::inserter(_pendingClassifiers, std::end(_pendingClassifiers)),
|
||||
[](auto& classifier) { return classifier.get(); });
|
||||
}
|
||||
|
||||
for (std::size_t i {}; i < classifiers.size(); ++i)
|
||||
loadClassifier(std::move(classifiers[i]), classifierTypes[i], forceReload, progressCallback);
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Recommendation engines loaded!";
|
||||
}
|
||||
|
||||
void
|
||||
Engine::setClassifierPriorities(const std::vector<ClassifierType>& classifierPriorities)
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
|
||||
_classifierPriorities.clear();
|
||||
std::transform(std::cbegin(classifierPriorities), std::cend(classifierPriorities), std::back_inserter(_classifierPriorities), [](std::string_view name) { return std::string {name}; });
|
||||
_classifierPriorities = classifierPriorities;
|
||||
}
|
||||
|
||||
void
|
||||
Engine::clearClassifiers()
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
std::unique_lock lock {_classifiersMutex};
|
||||
|
||||
_classifiers.clear();
|
||||
}
|
||||
|
||||
void
|
||||
Engine::initAndAddClassifier(std::unique_ptr<IClassifier> classifier, bool databaseChanged)
|
||||
Engine::loadClassifier(std::unique_ptr<IClassifier> classifier,
|
||||
ClassifierType classifierType,
|
||||
bool forceReload,
|
||||
const ProgressCallback& progressCallback)
|
||||
{
|
||||
PendingClassifierHandler pendingClassifier {*this, *classifier.get()};
|
||||
IClassifier* rawClassifier {classifier.get()};
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Initializing classifier '" << classifier->getName() << "'...";
|
||||
bool res {classifier->init(_dbSession, databaseChanged)};
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Initializing classifier '" << classifier->getName() << "': " << (res ? "SUCCESS" : "FAILURE");
|
||||
bool res {};
|
||||
if (!_loadCancelled)
|
||||
{
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Initializing classifier '" << classifier->getName() << "'...";
|
||||
|
||||
auto progress {[&](IClassifier::Progress progress)
|
||||
{
|
||||
progressCallback(Progress {progress.processedElems, progress.totalElems});
|
||||
}};
|
||||
|
||||
res = classifier->load(_db.getTLSSession(), forceReload, progressCallback ? progress : IClassifier::ProgressCallback {});
|
||||
|
||||
LMS_LOG(RECOMMENDATION, INFO) << "Initializing classifier '" << classifier->getName() << "': " << (res ? "SUCCESS" : "FAILURE");
|
||||
}
|
||||
|
||||
if (res)
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
std::unique_lock lock {_classifiersMutex};
|
||||
|
||||
_classifiers.emplace(classifier->getName(), std::move(classifier));
|
||||
_classifiers.emplace(classifierType, std::move(classifier));
|
||||
}
|
||||
|
||||
{
|
||||
std::scoped_lock lock {_controlMutex};
|
||||
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "About to erase. _pendingClassifiers size = " << _pendingClassifiers.size();
|
||||
_pendingClassifiers.erase(rawClassifier);
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Erased. _pendingClassifiers size = " << _pendingClassifiers.size();
|
||||
}
|
||||
|
||||
_pendingClassifiersCondvar.notify_one();
|
||||
|
||||
}
|
||||
|
||||
void
|
||||
Engine::cancelPendingClassifiers()
|
||||
Engine::cancelLoad()
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Cancelling loading...";
|
||||
|
||||
std::unique_lock lock {_controlMutex};
|
||||
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Still " << _pendingClassifiers.size() << " pending classifiers!";
|
||||
|
||||
_loadCancelled = true;
|
||||
|
||||
for (IClassifier* classifier : _pendingClassifiers)
|
||||
classifier->requestCancelInit();
|
||||
}
|
||||
classifier->requestCancelLoad();
|
||||
|
||||
void
|
||||
Engine::addPendingClassifier(IClassifier& classifier)
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
_pendingClassifiersCondvar.wait(lock, [this] {return _pendingClassifiers.empty();});
|
||||
_loadCancelled = false;
|
||||
|
||||
_pendingClassifiers.insert(&classifier);
|
||||
}
|
||||
|
||||
void
|
||||
Engine::removePendingClassifier(IClassifier& classifier)
|
||||
{
|
||||
std::unique_lock<std::shared_mutex> lock {_classifiersMutex};
|
||||
|
||||
_pendingClassifiers.erase(&classifier);
|
||||
LMS_LOG(RECOMMENDATION, DEBUG) << "Cancelling loading DONE";
|
||||
}
|
||||
|
||||
} // ns Similarity
|
||||
|
||||
Reference in New Issue
Block a user