Simplified the similarity engine loading (since there is no more the features-based engine)

This commit is contained in:
emeric
2023-11-10 20:59:25 +01:00
parent 8ea532fa26
commit 636b70450b
10 changed files with 504 additions and 711 deletions
@@ -33,240 +33,91 @@
namespace Recommendation namespace Recommendation
{ {
namespace
{
Database::ScanSettings::SimilarityEngineType getSimilarityEngineType(Database::Session& session)
{
auto transaction{ session.createSharedTransaction() };
static return Database::ScanSettings::get(session)->getSimilarityEngineType();
std::string_view }
engineTypeToString(EngineType engineType) }
{
switch (engineType)
{
case EngineType::Clusters: return "clusters";
case EngineType::Features: return "features";
}
throw LmsException {"Internal error"}; std::unique_ptr<IRecommendationService> createRecommendationService(Database::Db& db)
} {
return std::make_unique<RecommendationService>(db);
}
std::unique_ptr<IRecommendationService> RecommendationService::RecommendationService(Database::Db& db)
createRecommendationService(Database::Db& db) : _db{ db }
{ {
return std::make_unique<RecommendationService>(db); load();
} }
RecommendationService::RecommendationService(Database::Db& db) TrackContainer RecommendationService::findSimilarTracks(Database::TrackListId trackListId, std::size_t maxCount) const
: _db {db} {
{ TrackContainer res;
}
TrackContainer if (!_engine)
RecommendationService::findSimilarTracks(Database::TrackListId trackListId, std::size_t maxCount) const return res;
{
TrackContainer res;
std::shared_lock lock {_enginesMutex}; return _engine->findSimilarTracksFromTrackList(trackListId, maxCount);
for (const auto& engineType : _enginePriorities) }
{
auto itEngine {_engines.find(engineType)};
if (itEngine == std::cend(_engines))
continue;
res = itEngine->second->findSimilarTracksFromTrackList(trackListId, maxCount); TrackContainer RecommendationService::findSimilarTracks(const std::vector<Database::TrackId>& trackIds, std::size_t maxCount) const
if (!res.empty()) {
break; TrackContainer res;
}
return res; if (!_engine)
} return res;
TrackContainer return _engine->findSimilarTracks(trackIds, maxCount);
RecommendationService::findSimilarTracks(const std::vector<Database::TrackId>& trackIds, std::size_t maxCount) const }
{
TrackContainer res;
std::shared_lock lock {_enginesMutex}; ReleaseContainer RecommendationService::getSimilarReleases(Database::ReleaseId releaseId, std::size_t maxCount) const
for (EngineType engineType : _enginePriorities) {
{ ReleaseContainer res;
auto itEngine {_engines.find(engineType)};
if (itEngine == std::cend(_engines))
continue;
LMS_LOG(RECOMMENDATION, DEBUG) << "Trying engine '" << engineTypeToString(engineType) << "' to get similar tracks"; if (!_engine)
return res;
const IEngine& engine {*itEngine->second}; return _engine->getSimilarReleases(releaseId, maxCount);;
res = engine.findSimilarTracks(trackIds, maxCount); }
if (!res.empty())
{
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar tracks using engine '" << engineTypeToString(engineType) << "'";
break;
}
}
return res; ArtistContainer RecommendationService::getSimilarArtists(Database::ArtistId artistId, EnumSet<Database::TrackArtistLinkType> linkTypes, std::size_t maxCount) const
} {
ArtistContainer res;
ReleaseContainer if (!_engine)
RecommendationService::getSimilarReleases(Database::ReleaseId releaseId, std::size_t maxCount) const return res;
{
ReleaseContainer res;
std::shared_lock lock {_enginesMutex}; return _engine->getSimilarArtists(artistId, linkTypes, maxCount);
for (EngineType engineType : _enginePriorities)
{
auto itEngine {_engines.find(engineType)};
if (itEngine == std::cend(_engines))
continue;
LMS_LOG(RECOMMENDATION, DEBUG) << "Trying engine '" << engineTypeToString(engineType) << "' to get similar releases"; return res;
}
const IEngine& engine {*itEngine->second}; void RecommendationService::load()
res = engine.getSimilarReleases(releaseId, maxCount); {
if (!res.empty()) using namespace Database;
{
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar releases using engine '" << engineTypeToString(engineType) << "'";
break;
}
LMS_LOG(RECOMMENDATION, DEBUG) << "No result using engine '" << engineTypeToString(engineType) << "'"; switch (getSimilarityEngineType(_db.getTLSSession()))
} {
case ScanSettings::SimilarityEngineType::Clusters:
if (_engineType != EngineType::Clusters)
{
_engineType = EngineType::Clusters;
_engine = createClustersEngine(_db);
}
break;
return res; case ScanSettings::SimilarityEngineType::Features:
} case ScanSettings::SimilarityEngineType::None:
_engineType.reset();
ArtistContainer _engine.reset();
RecommendationService::getSimilarArtists(Database::ArtistId artistId, EnumSet<Database::TrackArtistLinkType> linkTypes, std::size_t maxCount) const break;
{ }
ArtistContainer res;
std::shared_lock lock {_enginesMutex};
for (EngineType engineType : _enginePriorities)
{
auto itEngine {_engines.find(engineType)};
if (itEngine == std::cend(_engines))
continue;
LMS_LOG(RECOMMENDATION, DEBUG) << "Trying engine '" << engineTypeToString(engineType) << "' to get similar artists";
const IEngine& engine {*itEngine->second};
res = engine.getSimilarArtists(artistId, linkTypes, maxCount);
if (!res.empty())
{
LMS_LOG(RECOMMENDATION, DEBUG) << "Got " << res.size() << " similar artists using engine '" << engineTypeToString(engineType) << "'";
return res;
}
}
return res;
}
static
Database::ScanSettings::SimilarityEngineType
getSimilarityEngineType(Database::Session& session)
{
auto transaction {session.createSharedTransaction()};
return Database::ScanSettings::get(session)->getSimilarityEngineType();
}
void
RecommendationService::load(bool forceReload, const ProgressCallback& progressCallback)
{
using namespace Database;
LMS_LOG(RECOMMENDATION, INFO) << "Reloading recommendation engines...";
EngineContainer enginesToLoad;
{
std::unique_lock controlLock {_controlMutex};
{
std::unique_lock lock {_enginesMutex};
_engines.clear();
}
switch (getSimilarityEngineType(_db.getTLSSession()))
{
case ScanSettings::SimilarityEngineType::Clusters:
_enginePriorities = {EngineType::Clusters};
enginesToLoad.try_emplace(EngineType::Clusters, createClustersEngine(_db));
break;
case ScanSettings::SimilarityEngineType::Features:
_enginePriorities = {EngineType::Features, EngineType::Clusters};
// not same order since clusters is faster to load
enginesToLoad.try_emplace(EngineType::Clusters, createClustersEngine(_db));
enginesToLoad.try_emplace(EngineType::Features, createFeaturesEngine(_db));
break;
case ScanSettings::SimilarityEngineType::None:
_enginePriorities.clear();
break;
}
assert(_pendingEngines.empty());
for (auto& [engineType, engine] : enginesToLoad)
_pendingEngines.push_back(engine.get());
}
for (auto& [engineType, engine] : enginesToLoad)
loadPendingEngine(engineType, std::move(engine), forceReload, progressCallback);
_pendingEnginesCondvar.notify_all();
LMS_LOG(RECOMMENDATION, INFO) << "Recommendation engines loaded!";
}
void
RecommendationService::loadPendingEngine(EngineType engineType, std::unique_ptr<IEngine> engine, bool forceReload, const ProgressCallback& progressCallback)
{
if (!_loadCancelled)
{
LMS_LOG(RECOMMENDATION, INFO) << "Initializing engine '" << engineTypeToString(engineType) << "'...";
auto progress {[&](const Progress& progress)
{
progressCallback(progress);
}};
engine->load(forceReload, progressCallback ? progress : ProgressCallback {});
LMS_LOG(RECOMMENDATION, INFO) << "Initializing engine '" << engineTypeToString(engineType) << "': " << (_loadCancelled ? "aborted" : "complete");
}
{
std::scoped_lock lock {_controlMutex};
_pendingEngines.erase(std::find(std::begin(_pendingEngines), std::end(_pendingEngines), engine.get()));
}
if (!_loadCancelled)
{
std::unique_lock lock {_enginesMutex};
_engines.emplace(engineType, std::move(engine));
}
}
void
RecommendationService::cancelLoad()
{
LMS_LOG(RECOMMENDATION, DEBUG) << "Cancelling loading...";
std::unique_lock controlLock {_controlMutex};
assert(!_loadCancelled);
_loadCancelled = true;
LMS_LOG(RECOMMENDATION, DEBUG) << "Still " << _pendingEngines.size() << " pending engines!";
for (IEngine* engine : _pendingEngines)
{
engine->requestCancelLoad();
}
_pendingEnginesCondvar.wait(controlLock, [this] {return _pendingEngines.empty();});
_loadCancelled = false;
LMS_LOG(RECOMMENDATION, DEBUG) << "Cancelling loading DONE";
}
if (_engine)
_engine->load(false);
}
} // ns Similarity } // ns Similarity
@@ -19,67 +19,49 @@
#pragma once #pragma once
#include <condition_variable> #include <optional>
#include <mutex>
#include <shared_mutex>
#include <unordered_map>
#include <vector>
#include "services/recommendation/IRecommendationService.hpp" #include "services/recommendation/IRecommendationService.hpp"
#include "IEngine.hpp" #include "IEngine.hpp"
namespace Database namespace Database
{ {
class Db; class Db;
} }
namespace Recommendation namespace Recommendation
{ {
enum class EngineType enum class EngineType
{ {
Clusters, Clusters,
Features, Features,
}; };
class RecommendationService : public IRecommendationService class RecommendationService : public IRecommendationService
{ {
public: public:
RecommendationService(Database::Db& db); RecommendationService(Database::Db& db);
~RecommendationService() = default; ~RecommendationService() = default;
RecommendationService(const RecommendationService&) = delete; RecommendationService(const RecommendationService&) = delete;
RecommendationService(RecommendationService&&) = delete; RecommendationService& operator=(const RecommendationService&) = delete;
RecommendationService& operator=(const RecommendationService&) = delete;
RecommendationService& operator=(RecommendationService&&) = delete;
private: private:
void load(bool forceReload, const ProgressCallback& progressCallback) override; void load() override;
void cancelLoad() override;
TrackContainer findSimilarTracks(Database::TrackListId tracklistId, std::size_t maxCount) const override; TrackContainer findSimilarTracks(Database::TrackListId tracklistId, std::size_t maxCount) const override;
TrackContainer findSimilarTracks(const std::vector<Database::TrackId>& tracksId, std::size_t maxCount) const override; TrackContainer findSimilarTracks(const std::vector<Database::TrackId>& tracksId, std::size_t maxCount) const override;
ReleaseContainer getSimilarReleases(Database::ReleaseId releaseId, std::size_t maxCount) const override; ReleaseContainer getSimilarReleases(Database::ReleaseId releaseId, std::size_t maxCount) const override;
ArtistContainer getSimilarArtists(Database::ArtistId artistId, EnumSet<Database::TrackArtistLinkType> linkTypes, std::size_t maxCount) const override; ArtistContainer getSimilarArtists(Database::ArtistId artistId, EnumSet<Database::TrackArtistLinkType> linkTypes, std::size_t maxCount) const override;
void setEnginePriorities(const std::vector<EngineType>& engineTypes); void setEnginePriorities(const std::vector<EngineType>& engineTypes);
void clearEngines(); void clearEngines();
void loadPendingEngine(EngineType engineType, std::unique_ptr<IEngine> engine, bool forceReload, const ProgressCallback& progressCallback); void loadPendingEngine(EngineType engineType, std::unique_ptr<IEngine> engine, bool forceReload, const ProgressCallback& progressCallback);
Database::Db& _db; Database::Db& _db;
std::optional<EngineType> _engineType;
std::mutex _controlMutex; std::unique_ptr<IEngine> _engine;
bool _loadCancelled {}; };
using EngineContainer = std::unordered_map<EngineType, std::unique_ptr<IEngine>>;
EngineContainer _engines;
mutable std::shared_mutex _enginesMutex;
std::vector<IEngine*> _pendingEngines;
std::shared_mutex _pendingEnginesMutex;
std::condition_variable _pendingEnginesCondvar;
std::vector<EngineType> _enginePriorities; // ordered by priority
};
} // ns Recommendation } // ns Recommendation
@@ -20,6 +20,7 @@
#pragma once #pragma once
#include <memory> #include <memory>
#include <vector>
#include "utils/EnumSet.hpp" #include "utils/EnumSet.hpp"
#include "services/database/TrackListId.hpp" #include "services/database/TrackListId.hpp"
#include "services/database/Types.hpp" #include "services/database/Types.hpp"
@@ -37,8 +38,7 @@ namespace Recommendation
public: public:
virtual ~IRecommendationService() = default; virtual ~IRecommendationService() = default;
virtual void load(bool forceReload, const ProgressCallback& progressCallback = {}) = 0; virtual void load() = 0;
virtual void cancelLoad() = 0; // wait for cancel done
virtual TrackContainer findSimilarTracks(Database::TrackListId tracklistId, std::size_t maxCount) const = 0; virtual TrackContainer findSimilarTracks(Database::TrackListId tracklistId, std::size_t maxCount) const = 0;
virtual TrackContainer findSimilarTracks(const std::vector<Database::TrackId>& tracksId, std::size_t maxCount) const = 0; virtual TrackContainer findSimilarTracks(const std::vector<Database::TrackId>& tracksId, std::size_t maxCount) const = 0;
+355 -382
View File
@@ -25,7 +25,6 @@
#include "services/database/Cluster.hpp" #include "services/database/Cluster.hpp"
#include "services/database/TrackFeatures.hpp" #include "services/database/TrackFeatures.hpp"
#include "services/database/ScanSettings.hpp" #include "services/database/ScanSettings.hpp"
#include "services/recommendation/IRecommendationService.hpp"
#include "utils/Exception.hpp" #include "utils/Exception.hpp"
#include "utils/IConfig.hpp" #include "utils/IConfig.hpp"
#include "utils/Logger.hpp" #include "utils/Logger.hpp"
@@ -38,387 +37,361 @@
#include "ScanStepScanFiles.hpp" #include "ScanStepScanFiles.hpp"
#include "ScanStepComputeClusterStats.hpp" #include "ScanStepComputeClusterStats.hpp"
using namespace Database; namespace Scanner
namespace {
Wt::WDate
getNextMonday(Wt::WDate current)
{ {
do using namespace Database;
{
current = current.addDays(1); namespace
} while (current.dayOfWeek() != 1); {
Wt::WDate getNextMonday(Wt::WDate current)
return current; {
} do
{
Wt::WDate current = current.addDays(1);
getNextFirstOfMonth(Wt::WDate current) } while (current.dayOfWeek() != 1);
{
do return current;
{ }
current = current.addDays(1);
} while (current.day() != 1); Wt::WDate getNextFirstOfMonth(Wt::WDate current)
{
return current; do
} {
current = current.addDays(1);
} // namespace } while (current.day() != 1);
namespace Scanner { return current;
}
std::unique_ptr<IScannerService> } // namespace
createScannerService(Db& db, Recommendation::IRecommendationService& recommendationService)
{ std::unique_ptr<IScannerService> createScannerService(Db& db)
return std::make_unique<ScannerService>(db, recommendationService); {
} return std::make_unique<ScannerService>(db);
}
ScannerService::ScannerService(Db& db, Recommendation::IRecommendationService& recommendationService)
: _recommendationService {recommendationService} ScannerService::ScannerService(Db& db)
, _db {db} : _db{ db }
, _dbSession {db} , _dbSession{ db }
{ {
_ioService.setThreadCount(1); _ioService.setThreadCount(1);
refreshScanSettings(); refreshScanSettings();
start(); start();
} }
ScannerService::~ScannerService() ScannerService::~ScannerService()
{ {
LMS_LOG(DBUPDATER, INFO) << "Stopping service..."; LMS_LOG(DBUPDATER, INFO) << "Stopping service...";
stop(); stop();
LMS_LOG(DBUPDATER, INFO) << "Service stopped!"; LMS_LOG(DBUPDATER, INFO) << "Service stopped!";
} }
void void ScannerService::start()
ScannerService::start() {
{ std::scoped_lock lock{ _controlMutex };
std::scoped_lock lock {_controlMutex};
_ioService.post([this]
_ioService.post([this] {
{ if (_abortScan)
if (_abortScan) return;
return;
scheduleNextScan();
_recommendationService.load(false, });
[](const Recommendation::Progress& progress)
{ _ioService.start();
LMS_LOG(DBUPDATER, DEBUG) << "Reloading recommendation : " << progress.processedElems << "/" << progress.totalElems; }
});
scheduleNextScan(); void ScannerService::stop()
}); {
std::scoped_lock lock{ _controlMutex };
_ioService.start();
} _abortScan = true;
_scheduleTimer.cancel();
void _ioService.stop();
ScannerService::stop() }
{
std::scoped_lock lock {_controlMutex}; void ScannerService::abortScan()
{
_abortScan = true; LMS_LOG(DBUPDATER, DEBUG) << "Aborting scan...";
_scheduleTimer.cancel(); std::scoped_lock lock{ _controlMutex };
_recommendationService.cancelLoad();
_ioService.stop(); LMS_LOG(DBUPDATER, DEBUG) << "Waiting for the scan to abort...";
}
_abortScan = true;
void _scheduleTimer.cancel();
ScannerService::abortScan() _ioService.stop();
{ LMS_LOG(DBUPDATER, DEBUG) << "Scan abort done!";
LMS_LOG(DBUPDATER, DEBUG) << "Aborting scan...";
std::scoped_lock lock {_controlMutex}; _abortScan = false;
_ioService.start();
LMS_LOG(DBUPDATER, DEBUG) << "Waiting for the scan to abort..."; }
_abortScan = true; void ScannerService::requestImmediateScan(bool force)
_scheduleTimer.cancel(); {
_recommendationService.cancelLoad(); abortScan();
_ioService.stop(); _ioService.post([=]()
LMS_LOG(DBUPDATER, DEBUG) << "Scan abort done!"; {
if (_abortScan)
_abortScan = false; return;
_ioService.start();
} scheduleScan(force);
});
void }
ScannerService::requestImmediateScan(bool force)
{ void ScannerService::requestReload()
abortScan(); {
_ioService.post([=]() abortScan();
{ _ioService.post([=]()
if (_abortScan) {
return; if (_abortScan)
return;
scheduleScan(force);
}); scheduleNextScan();
} });
}
void
ScannerService::requestReload() ScannerService::Status ScannerService::getStatus() const
{ {
abortScan(); Status res;
_ioService.post([=]()
{ std::shared_lock lock{ _statusMutex };
if (_abortScan)
return; res.currentState = _curState;
res.nextScheduledScan = _nextScheduledScan;
scheduleNextScan(); res.lastCompleteScanStats = _lastCompleteScanStats;
}); res.currentScanStepStats = _currentScanStepStats;
}
return res;
ScannerService::Status }
ScannerService::getStatus() const
{ void ScannerService::scheduleNextScan()
Status res; {
LMS_LOG(DBUPDATER, DEBUG) << "Scheduling next scan";
std::shared_lock lock {_statusMutex};
refreshScanSettings();
res.currentState = _curState;
res.nextScheduledScan = _nextScheduledScan; const Wt::WDateTime now{ Wt::WDateTime::currentDateTime() };
res.lastCompleteScanStats = _lastCompleteScanStats;
res.currentScanStepStats = _currentScanStepStats; Wt::WDateTime nextScanDateTime;
switch (_settings.updatePeriod)
return res; {
} case ScanSettings::UpdatePeriod::Daily:
if (now.time() < _settings.startTime)
void nextScanDateTime = { now.date(), _settings.startTime };
ScannerService::scheduleNextScan() else
{ nextScanDateTime = { now.date().addDays(1), _settings.startTime };
LMS_LOG(DBUPDATER, DEBUG) << "Scheduling next scan"; break;
refreshScanSettings(); case ScanSettings::UpdatePeriod::Weekly:
if (now.time() < _settings.startTime && now.date().dayOfWeek() == 1)
const Wt::WDateTime now {Wt::WDateTime::currentDateTime()}; nextScanDateTime = { now.date(), _settings.startTime };
else
Wt::WDateTime nextScanDateTime; nextScanDateTime = { getNextMonday(now.date()), _settings.startTime };
switch (_settings.updatePeriod) break;
{
case ScanSettings::UpdatePeriod::Daily: case ScanSettings::UpdatePeriod::Monthly:
if (now.time() < _settings.startTime) if (now.time() < _settings.startTime && now.date().day() == 1)
nextScanDateTime = {now.date(), _settings.startTime}; nextScanDateTime = { now.date(), _settings.startTime };
else else
nextScanDateTime = {now.date().addDays(1), _settings.startTime}; nextScanDateTime = { getNextFirstOfMonth(now.date()), _settings.startTime };
break; break;
case ScanSettings::UpdatePeriod::Weekly: case ScanSettings::UpdatePeriod::Hourly:
if (now.time() < _settings.startTime && now.date().dayOfWeek() == 1) nextScanDateTime = { now.date(), now.time().addSecs(3600) };
nextScanDateTime = {now.date(), _settings.startTime}; break;
else
nextScanDateTime = {getNextMonday(now.date()), _settings.startTime}; case ScanSettings::UpdatePeriod::Never:
break; LMS_LOG(DBUPDATER, INFO) << "Auto scan disabled!";
break;
case ScanSettings::UpdatePeriod::Monthly: }
if (now.time() < _settings.startTime && now.date().day() == 1)
nextScanDateTime = {now.date(), _settings.startTime}; if (nextScanDateTime.isValid())
else scheduleScan(false, nextScanDateTime);
nextScanDateTime = {getNextFirstOfMonth(now.date()), _settings.startTime};
break; {
std::unique_lock lock{ _statusMutex };
case ScanSettings::UpdatePeriod::Hourly: _curState = nextScanDateTime.isValid() ? State::Scheduled : State::NotScheduled;
nextScanDateTime = {now.date(), now.time().addSecs(3600)}; _nextScheduledScan = nextScanDateTime;
break; }
case ScanSettings::UpdatePeriod::Never: _events.scanScheduled.emit(_nextScheduledScan);
LMS_LOG(DBUPDATER, INFO) << "Auto scan disabled!"; }
break;
} void ScannerService::scheduleScan(bool force, const Wt::WDateTime& dateTime)
{
if (nextScanDateTime.isValid()) auto cb{ [=](boost::system::error_code ec)
scheduleScan(false, nextScanDateTime); {
if (ec)
{ return;
std::unique_lock lock {_statusMutex};
_curState = nextScanDateTime.isValid() ? State::Scheduled : State::NotScheduled; scan(force);
_nextScheduledScan = nextScanDateTime; } };
}
if (dateTime.isNull())
_events.scanScheduled.emit(_nextScheduledScan); {
} LMS_LOG(DBUPDATER, INFO) << "Scheduling next scan right now";
_scheduleTimer.expires_from_now(std::chrono::seconds{ 0 });
void _scheduleTimer.async_wait(cb);
ScannerService::scheduleScan(bool force, const Wt::WDateTime& dateTime) }
{ else
auto cb {[=](boost::system::error_code ec) {
{ std::chrono::system_clock::time_point timePoint{ dateTime.toTimePoint() };
if (ec) std::time_t t{ std::chrono::system_clock::to_time_t(timePoint) };
return; char ctimeStr[26];
scan(force); LMS_LOG(DBUPDATER, INFO) << "Scheduling next scan at " << std::string(::ctime_r(&t, ctimeStr));
}}; _scheduleTimer.expires_at(timePoint);
_scheduleTimer.async_wait(cb);
if (dateTime.isNull()) }
{ }
LMS_LOG(DBUPDATER, INFO) << "Scheduling next scan right now";
_scheduleTimer.expires_from_now(std::chrono::seconds {0}); void ScannerService::scan(bool forceScan)
_scheduleTimer.async_wait(cb); {
} _events.scanStarted.emit();
else
{ {
std::chrono::system_clock::time_point timePoint {dateTime.toTimePoint()}; std::unique_lock lock{ _statusMutex };
std::time_t t {std::chrono::system_clock::to_time_t(timePoint)}; _curState = State::InProgress;
char ctimeStr[26]; _nextScheduledScan = {};
}
LMS_LOG(DBUPDATER, INFO) << "Scheduling next scan at " << std::string(::ctime_r(&t, ctimeStr));
_scheduleTimer.expires_at(timePoint);
_scheduleTimer.async_wait(cb); LMS_LOG(UI, INFO) << "New scan started!";
}
} refreshScanSettings();
void IScanStep::ScanContext scanContext{ _settings.mediaDirectory, forceScan, ScanStats {}, ScanStepStats {} };
ScannerService::scan(bool forceScan) ScanStats& stats{ scanContext.stats };
{ stats.startTime = Wt::WDateTime::currentDateTime();
_events.scanStarted.emit();
for (auto& scanStep : _scanSteps)
{ {
std::unique_lock lock {_statusMutex}; LMS_LOG(DBUPDATER, DEBUG) << "Starting scan step '" << scanStep->getStepName() << "'";
_curState = State::InProgress; scanContext.currentStepStats = ScanStepStats{ Wt::WDateTime::currentDateTime(), scanStep->getStep() };
_nextScheduledScan = {};
} notifyInProgress(scanContext.currentStepStats);
scanStep->process(scanContext);
notifyInProgress(scanContext.currentStepStats);
LMS_LOG(UI, INFO) << "New scan started!"; LMS_LOG(DBUPDATER, DEBUG) << "Completed scan step '" << scanStep->getStepName() << "'";
}
refreshScanSettings();
LMS_LOG(DBUPDATER, INFO) << "Scan " << (_abortScan ? "aborted" : "complete") << ". Changes = " << stats.nbChanges() << " (added = " << stats.additions << ", removed = " << stats.deletions << ", updated = " << stats.updates << "), Not changed = " << stats.skips << ", Scanned = " << stats.scans << " (errors = " << stats.errors.size() << "), features fetched = " << stats.featuresFetched << ", duplicates = " << stats.duplicates.size();
IScanStep::ScanContext scanContext {_settings.mediaDirectory, forceScan, ScanStats {}, ScanStepStats {}};
ScanStats& stats {scanContext.stats}; _dbSession.analyze();
stats.startTime = Wt::WDateTime::currentDateTime();
if (!_abortScan)
for (auto& scanStep : _scanSteps) {
{ stats.stopTime = Wt::WDateTime::currentDateTime();
LMS_LOG(DBUPDATER, DEBUG) << "Starting scan step '" << scanStep->getStepName() << "'"; {
scanContext.currentStepStats = ScanStepStats {Wt::WDateTime::currentDateTime(), scanStep->getStep()}; std::unique_lock lock{ _statusMutex };
notifyInProgress(scanContext.currentStepStats); _lastCompleteScanStats = stats;
scanStep->process(scanContext); _currentScanStepStats.reset();
notifyInProgress(scanContext.currentStepStats); }
LMS_LOG(DBUPDATER, DEBUG) << "Completed scan step '" << scanStep->getStepName() << "'";
} LMS_LOG(DBUPDATER, DEBUG) << "Scan not aborted, scheduling next scan!";
scheduleNextScan();
LMS_LOG(DBUPDATER, INFO) << "Scan " << (_abortScan ? "aborted" : "complete") << ". Changes = " << stats.nbChanges() << " (added = " << stats.additions << ", removed = " << stats.deletions << ", updated = " << stats.updates << "), Not changed = " << stats.skips << ", Scanned = " << stats.scans << " (errors = " << stats.errors.size() << "), features fetched = " << stats.featuresFetched << ", duplicates = " << stats.duplicates.size();
_events.scanComplete.emit(stats);
_dbSession.analyze(); }
else
if (!_abortScan) {
{ LMS_LOG(DBUPDATER, DEBUG) << "Scan aborted, not scheduling next scan!";
stats.stopTime = Wt::WDateTime::currentDateTime();
{ std::unique_lock lock{ _statusMutex };
std::unique_lock lock {_statusMutex};
_curState = State::NotScheduled;
_lastCompleteScanStats = stats; _currentScanStepStats.reset();
_currentScanStepStats.reset(); }
} }
LMS_LOG(DBUPDATER, DEBUG) << "Scan not aborted, scheduling next scan!"; void ScannerService::refreshScanSettings()
scheduleNextScan(); {
ScannerSettings newSettings{ readSettings() };
_events.scanComplete.emit(stats); if (_settings == newSettings)
} return;
else
{ LMS_LOG(DBUPDATER, DEBUG) << "Scanner settings updated";
LMS_LOG(DBUPDATER, DEBUG) << "Scan aborted, not scheduling next scan!"; LMS_LOG(DBUPDATER, DEBUG) << "skipDuplicateMBID = " << newSettings.skipDuplicateMBID;
LMS_LOG(DBUPDATER, DEBUG) << "Using scan settings version " << newSettings.scanVersion;
std::unique_lock lock {_statusMutex};
_settings = std::move(newSettings);
_curState = State::NotScheduled;
_currentScanStepStats.reset(); auto cbFunc{ [this](const ScanStepStats& stats)
} {
} notifyInProgressIfNeeded(stats);
} };
void
ScannerService::refreshScanSettings() ScanStepBase::InitParams params
{ {
ScannerSettings newSettings {readSettings()}; _settings,
if (_settings == newSettings) cbFunc,
return; _abortScan,
_db
LMS_LOG(DBUPDATER, DEBUG) << "Scanner settings updated"; };
LMS_LOG(DBUPDATER, DEBUG) << "skipDuplicateMBID = " << newSettings.skipDuplicateMBID;
LMS_LOG(DBUPDATER, DEBUG) << "Using scan settings version " << newSettings.scanVersion; _scanSteps.clear();
_scanSteps.push_back(std::make_unique<ScanStepDiscoverFiles>(params));
_settings = std::move(newSettings); _scanSteps.push_back(std::make_unique<ScanStepScanFiles>(params));
_scanSteps.push_back(std::make_unique<ScanStepRemoveOrphanDbFiles>(params));
auto cbFunc {[this](const ScanStepStats& stats) _scanSteps.push_back(std::make_unique<ScanStepComputeClusterStats>(params));
{ _scanSteps.push_back(std::make_unique<ScanStepCheckDuplicatedDbFiles>(params));
notifyInProgressIfNeeded(stats); }
}};
ScannerSettings ScannerService::readSettings()
ScanStepBase::InitParams params {
{ ScannerSettings newSettings;
_settings,
cbFunc, newSettings.skipDuplicateMBID = Service<IConfig>::get()->getBool("scanner-skip-duplicate-mbid", false);
_abortScan, {
_db auto transaction{ _dbSession.createSharedTransaction() };
};
const ScanSettings::pointer scanSettings{ ScanSettings::get(_dbSession) };
_scanSteps.clear();
_scanSteps.push_back(std::make_unique<ScanStepDiscoverFiles>(params)); newSettings.scanVersion = scanSettings->getScanVersion();
_scanSteps.push_back(std::make_unique<ScanStepScanFiles>(params)); newSettings.startTime = scanSettings->getUpdateStartTime();
_scanSteps.push_back(std::make_unique<ScanStepRemoveOrphanDbFiles>(params)); newSettings.updatePeriod = scanSettings->getUpdatePeriod();
_scanSteps.push_back(std::make_unique<ScanStepComputeClusterStats>(params));
_scanSteps.push_back(std::make_unique<ScanStepCheckDuplicatedDbFiles>(params)); {
} const auto fileExtensions{ scanSettings->getAudioFileExtensions() };
newSettings.supportedExtensions.reserve(fileExtensions.size());
ScannerSettings std::transform(std::cbegin(fileExtensions), std::end(fileExtensions), std::back_inserter(newSettings.supportedExtensions),
ScannerService::readSettings() [](const std::filesystem::path& extension) { return std::filesystem::path{ StringUtils::stringToLower(extension.string()) }; });
{ }
ScannerSettings newSettings; newSettings.mediaDirectory = scanSettings->getMediaDirectory();
newSettings.skipDuplicateMBID = Service<IConfig>::get()->getBool("scanner-skip-duplicate-mbid", false); const auto clusterTypes = scanSettings->getClusterTypes();
{ std::set<std::string> clusterTypeNames;
auto transaction {_dbSession.createSharedTransaction()};
std::transform(std::cbegin(clusterTypes), std::cend(clusterTypes),
const ScanSettings::pointer scanSettings {ScanSettings::get(_dbSession)}; std::inserter(clusterTypeNames, clusterTypeNames.begin()),
[](ClusterType::pointer clusterType) { return clusterType->getName(); });
newSettings.scanVersion = scanSettings->getScanVersion();
newSettings.startTime = scanSettings->getUpdateStartTime(); newSettings.clusterTypeNames = std::move(clusterTypeNames);
newSettings.updatePeriod = scanSettings->getUpdatePeriod(); }
{ return newSettings;
const auto fileExtensions {scanSettings->getAudioFileExtensions()}; }
newSettings.supportedExtensions.reserve(fileExtensions.size());
std::transform(std::cbegin(fileExtensions), std::end(fileExtensions), std::back_inserter(newSettings.supportedExtensions), void ScannerService::notifyInProgress(const ScanStepStats& stepStats)
[](const std::filesystem::path& extension) { return std::filesystem::path{ StringUtils::stringToLower(extension.string()) }; }); {
} {
newSettings.similarityServiceType = scanSettings->getSimilarityEngineType(); std::unique_lock lock{ _statusMutex };
newSettings.mediaDirectory = scanSettings->getMediaDirectory(); _currentScanStepStats = stepStats;
}
const auto clusterTypes = scanSettings->getClusterTypes();
std::set<std::string> clusterTypeNames; const std::chrono::system_clock::time_point now{ std::chrono::system_clock::now() };
_events.scanInProgress(stepStats);
std::transform(std::cbegin(clusterTypes), std::cend(clusterTypes), _lastScanInProgressEmit = now;
std::inserter(clusterTypeNames, clusterTypeNames.begin()), }
[](ClusterType::pointer clusterType) { return clusterType->getName(); });
void ScannerService::notifyInProgressIfNeeded(const ScanStepStats& stepStats)
newSettings.clusterTypeNames = std::move(clusterTypeNames); {
} std::chrono::system_clock::time_point now{ std::chrono::system_clock::now() };
return newSettings; if (std::chrono::duration_cast<std::chrono::seconds>(now - _lastScanInProgressEmit).count() > 1)
} notifyInProgress(stepStats);
}
void
ScannerService::notifyInProgress(const ScanStepStats& stepStats)
{
{
std::unique_lock lock {_statusMutex};
_currentScanStepStats = stepStats;
}
const std::chrono::system_clock::time_point now {std::chrono::system_clock::now()};
_events.scanInProgress(stepStats);
_lastScanInProgressEmit = now;
}
void
ScannerService::notifyInProgressIfNeeded(const ScanStepStats& stepStats)
{
std::chrono::system_clock::time_point now {std::chrono::system_clock::now()};
if (std::chrono::duration_cast<std::chrono::seconds>(now - _lastScanInProgressEmit).count() > 1)
notifyInProgress(stepStats);
}
} // namespace Scanner } // namespace Scanner
@@ -38,73 +38,65 @@
#include "IScanStep.hpp" #include "IScanStep.hpp"
#include "ScannerSettings.hpp" #include "ScannerSettings.hpp"
namespace Recommendation
{
class IRecommendationService;
}
namespace Scanner namespace Scanner
{ {
class ScannerService : public IScannerService class ScannerService : public IScannerService
{ {
public: public:
ScannerService(Database::Db& db, Recommendation::IRecommendationService& recommendationService); ScannerService(Database::Db& db);
~ScannerService(); ~ScannerService();
ScannerService(const ScannerService&) = delete; ScannerService(const ScannerService&) = delete;
ScannerService(ScannerService&&) = delete; ScannerService& operator=(const ScannerService&) = delete;
ScannerService& operator=(const ScannerService&) = delete;
ScannerService& operator=(ScannerService&&) = delete;
void requestReload() override; void requestReload() override;
void requestImmediateScan(bool force) override; void requestImmediateScan(bool force) override;
Status getStatus() const override; Status getStatus() const override;
Events& getEvents() override { return _events; } Events& getEvents() override { return _events; }
private: private:
void start(); void start();
void stop(); void stop();
// Job handling // Job handling
void scheduleNextScan(); void scheduleNextScan();
void scheduleScan(bool force, const Wt::WDateTime& dateTime = {}); void scheduleScan(bool force, const Wt::WDateTime& dateTime = {});
void abortScan(); void abortScan();
// Update database (scheduled callback) // Update database (scheduled callback)
void scan(bool force); void scan(bool force);
void scanMediaDirectory( const std::filesystem::path& mediaDirectory, bool forceScan, ScanStats& stats); void scanMediaDirectory(const std::filesystem::path& mediaDirectory, bool forceScan, ScanStats& stats);
// Helpers // Helpers
void refreshScanSettings(); void refreshScanSettings();
ScannerSettings readSettings(); ScannerSettings readSettings();
void reloadRecommendationService();
void notifyInProgressIfNeeded(const ScanStepStats& stats); void notifyInProgressIfNeeded(const ScanStepStats& stats);
void notifyInProgress(const ScanStepStats& stats); void notifyInProgress(const ScanStepStats& stats);
void reloadSimilarityEngine(ScanStats& stats); void reloadSimilarityEngine(ScanStats& stats);
Recommendation::IRecommendationService& _recommendationService; std::vector<std::unique_ptr<IScanStep>> _scanSteps;
std::vector<std::unique_ptr<IScanStep>> _scanSteps; std::mutex _controlMutex;
bool _abortScan{};
Wt::WIOService _ioService;
boost::asio::system_timer _scheduleTimer{ _ioService };
Events _events;
std::chrono::system_clock::time_point _lastScanInProgressEmit{};
Database::Db& _db;
Database::Session _dbSession;
std::mutex _controlMutex; mutable std::shared_mutex _statusMutex;
bool _abortScan {}; State _curState{ State::NotScheduled };
Wt::WIOService _ioService; std::optional<ScanStats> _lastCompleteScanStats;
boost::asio::system_timer _scheduleTimer {_ioService}; std::optional<ScanStepStats> _currentScanStepStats;
Events _events; Wt::WDateTime _nextScheduledScan;
std::chrono::system_clock::time_point _lastScanInProgressEmit {};
Database::Db& _db;
Database::Session _dbSession;
mutable std::shared_mutex _statusMutex; ScannerSettings _settings;
State _curState {State::NotScheduled}; };
std::optional<ScanStats> _lastCompleteScanStats;
std::optional<ScanStepStats> _currentScanStepStats;
Wt::WDateTime _nextScheduledScan;
ScannerSettings _settings;
};
} // Scanner } // Scanner
@@ -34,7 +34,6 @@ namespace Scanner
Wt::WTime startTime; Wt::WTime startTime;
Database::ScanSettings::UpdatePeriod updatePeriod {Database::ScanSettings::UpdatePeriod::Never}; Database::ScanSettings::UpdatePeriod updatePeriod {Database::ScanSettings::UpdatePeriod::Never};
std::vector<std::filesystem::path> supportedExtensions; std::vector<std::filesystem::path> supportedExtensions;
Database::ScanSettings::SimilarityEngineType similarityServiceType;
std::filesystem::path mediaDirectory; std::filesystem::path mediaDirectory;
bool skipDuplicateMBID {}; bool skipDuplicateMBID {};
std::set<std::string> clusterTypeNames; std::set<std::string> clusterTypeNames;
@@ -45,7 +44,6 @@ namespace Scanner
&& startTime == rhs.startTime && startTime == rhs.startTime
&& updatePeriod == rhs.updatePeriod && updatePeriod == rhs.updatePeriod
&& supportedExtensions == rhs.supportedExtensions && supportedExtensions == rhs.supportedExtensions
&& similarityServiceType == rhs.similarityServiceType
&& mediaDirectory == rhs.mediaDirectory && mediaDirectory == rhs.mediaDirectory
&& skipDuplicateMBID == rhs.skipDuplicateMBID && skipDuplicateMBID == rhs.skipDuplicateMBID
&& clusterTypeNames == rhs.clusterTypeNames; && clusterTypeNames == rhs.clusterTypeNames;
@@ -29,11 +29,6 @@ namespace Database
class Db; class Db;
} }
namespace Recommendation
{
class IRecommendationService;
}
namespace Scanner namespace Scanner
{ {
@@ -66,7 +61,7 @@ namespace Scanner
virtual Events& getEvents() = 0; virtual Events& getEvents() = 0;
}; };
std::unique_ptr<IScannerService> createScannerService(Database::Db& db, Recommendation::IRecommendationService& recommendationEngine); std::unique_ptr<IScannerService> createScannerService(Database::Db& db);
} // Scanner } // Scanner
+1 -1
View File
@@ -270,7 +270,7 @@ int main(int argc, char* argv[])
Service<Cover::ICoverService> coverService{ Cover::createCoverService(database, argv[0], server.appRoot() + "/images/unknown-cover.jpg") }; Service<Cover::ICoverService> coverService{ Cover::createCoverService(database, argv[0], server.appRoot() + "/images/unknown-cover.jpg") };
Service<Recommendation::IRecommendationService> recommendationService{ Recommendation::createRecommendationService(database) }; Service<Recommendation::IRecommendationService> recommendationService{ Recommendation::createRecommendationService(database) };
Service<Recommendation::IPlaylistGeneratorService> playlistGeneratorService{ Recommendation::createPlaylistGeneratorService(database, *recommendationService.get()) }; Service<Recommendation::IPlaylistGeneratorService> playlistGeneratorService{ Recommendation::createPlaylistGeneratorService(database, *recommendationService.get()) };
Service<Scanner::IScannerService> scannerService{ Scanner::createScannerService(database, *recommendationService) }; Service<Scanner::IScannerService> scannerService{ Scanner::createScannerService(database) };
scannerService->getEvents().scanComplete.connect([&] scannerService->getEvents().scanComplete.connect([&]
{ {
@@ -29,6 +29,7 @@
#include "services/database/Cluster.hpp" #include "services/database/Cluster.hpp"
#include "services/database/ScanSettings.hpp" #include "services/database/ScanSettings.hpp"
#include "services/database/Session.hpp" #include "services/database/Session.hpp"
#include "services/recommendation/IRecommendationService.hpp"
#include "services/scanner/IScannerService.hpp" #include "services/scanner/IScannerService.hpp"
#include "utils/Logger.hpp" #include "utils/Logger.hpp"
#include "utils/Service.hpp" #include "utils/Service.hpp"
@@ -232,6 +233,7 @@ DatabaseSettingsView::refreshView()
{ {
model->saveData(); model->saveData();
Service<Recommendation::IRecommendationService>::get()->load();
Service<Scanner::IScannerService>::get()->requestImmediateScan(false); Service<Scanner::IScannerService>::get()->requestImmediateScan(false);
LmsApp->notifyMsg(Notification::Type::Info, Wt::WString::tr("Lms.Admin.Database.database"), Wt::WString::tr("Lms.Admin.Database.settings-saved")); LmsApp->notifyMsg(Notification::Type::Info, Wt::WString::tr("Lms.Admin.Database.database"), Wt::WString::tr("Lms.Admin.Database.settings-saved"));
} }
@@ -157,16 +157,16 @@ int main(int argc, char* argv[])
Db db{ config->getPath("working-dir") / "lms.db" }; Db db{ config->getPath("working-dir") / "lms.db" };
Session session{ db }; Session session{ db };
std::cout << "Creating recommendation recommendationService..." << std::endl; std::cout << "Creating recommendation service..." << std::endl;
const auto recommendationService{ Recommendation::createRecommendationService(db) }; const auto recommendationService{ Recommendation::createRecommendationService(db) };
std::cout << "Recommendation recommendationService created!" << std::endl; std::cout << "Recommendation service created!" << std::endl;
std::cout << "Loading recommendation recommendationService..." << std::endl; std::cout << "Loading recommendation service..." << std::endl;
recommendationService->load(false); recommendationService->load();
unsigned maxSimilarityCount{ vm["max"].as<unsigned>() }; unsigned maxSimilarityCount{ vm["max"].as<unsigned>() };
std::cout << "Recommendation recommendationService loaded!" << std::endl; std::cout << "Recommendation service loaded!" << std::endl;
if (vm.count("tracks")) if (vm.count("tracks"))
dumpTracksRecommendation(db, *recommendationService, maxSimilarityCount); dumpTracksRecommendation(db, *recommendationService, maxSimilarityCount);