Send pending listens that have not been successfully sent

This commit is contained in:
emeric
2022-01-30 14:18:22 +01:00
parent 97c2571a54
commit 4177354afe
18 changed files with 232 additions and 154 deletions
@@ -105,7 +105,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getRecentArtists(session, userId, *scrobbler, clusterIds, linkType, range);
res = Database::Listen::getRecentArtists(session, userId, *scrobbler, clusterIds, linkType, range);
return res;
}
ScrobblingService::ReleaseContainer
@@ -120,7 +121,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getRecentReleases(session, userId, *scrobbler, clusterIds, range);
res = Database::Listen::getRecentReleases(session, userId, *scrobbler, clusterIds, range);
return res;
}
ScrobblingService::TrackContainer
@@ -135,7 +137,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getRecentTracks(session, userId, *scrobbler, clusterIds, range);
res = Database::Listen::getRecentTracks(session, userId, *scrobbler, clusterIds, range);
return res;
}
// Top
@@ -151,7 +154,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getTopArtists(session, userId, *scrobbler, clusterIds, linkType, range);
res = Database::Listen::getTopArtists(session, userId, *scrobbler, clusterIds, linkType, range);
return res;
}
ScrobblingService::ReleaseContainer
@@ -166,7 +170,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getTopReleases(session, userId, *scrobbler, clusterIds, range);
res = Database::Listen::getTopReleases(session, userId, *scrobbler, clusterIds, range);
return res;
}
ScrobblingService::TrackContainer
@@ -181,7 +186,8 @@ namespace Scrobbling
Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
return Database::Listen::getTopTracks(session, userId, *scrobbler, clusterIds, range);
res = Database::Listen::getTopTracks(session, userId, *scrobbler, clusterIds, range);
return res;
}
void
@@ -69,7 +69,7 @@ namespace Scrobbling::ListenBrainz
void
Scrobbler::listenStarted(const Listen& listen)
{
enqueListen(listen, Wt::WDateTime {});
_listensSynchronizer.enqueListenNow(listen);
}
void
@@ -78,23 +78,14 @@ namespace Scrobbling::ListenBrainz
if (duration && !canBeScrobbled(_db.getTLSSession(), listen.trackId, *duration))
return;
const Listen timedListen {listen};
const Wt::WDateTime now {Wt::WDateTime::currentDateTime()};
enqueListen(timedListen, now);
const TimedListen timedListen {listen, Wt::WDateTime::currentDateTime()};
_listensSynchronizer.enqueListen(timedListen);
}
void
Scrobbler::addTimedListen(const TimedListen& listen)
Scrobbler::addTimedListen(const TimedListen& timedListen)
{
assert(listen.listenedAt.isValid());
enqueListen(listen, listen.listenedAt);
}
void
Scrobbler::enqueListen(const Listen& listen, const Wt::WDateTime& timePoint)
{
_listensSynchronizer.enqueListen(listen, timePoint);
_listensSynchronizer.enqueListen(timedListen);
}
} // namespace Scrobbling::ListenBrainz
@@ -50,9 +50,6 @@ namespace Scrobbling::ListenBrainz
void listenFinished(const Listen& listen, std::optional<std::chrono::seconds> duration) override;
void addTimedListen(const TimedListen& listen) override;
// Submit listens
void enqueListen(const Listen& listen, const Wt::WDateTime& timePoint);
boost::asio::io_context& _ioContext;
Database::Db& _db;
std::string _baseAPIUrl;
@@ -42,6 +42,7 @@
#include "Utils.hpp"
#define LOG(sev) LMS_LOG(SCROBBLING, sev) << "[listenbrainz Synchronizer] - "
#define LOG_EX(sev) LMS_LOG_EX(Module::SCROBBLING, sev) << "[listenbrainz Synchronizer] - "
namespace
{
@@ -173,6 +174,8 @@ namespace
{
using namespace Database;
//LOG(DEBUG) << "Trying to match track' " << Wt::Json::serialize(metadata) << "'";
// first try to get the associated track using MBIDs, and then fallback on names
if (metadata.type("additional_info") == Wt::Json::Type::Object)
{
@@ -302,7 +305,20 @@ namespace Scrobbling::ListenBrainz
{
LOG(INFO) << "Starting Listens synchronizer, maxSyncListenCount = " << _maxSyncListenCount << ", _syncListensPeriod = " << _syncListensPeriod.count() << " hours";
scheduleGetListens(std::chrono::seconds {30});
scheduleSync(std::chrono::seconds {30});
}
void
ListensSynchronizer::enqueListen(const TimedListen& listen)
{
assert(listen.listenedAt.isValid());
enqueListen(listen, listen.listenedAt);
}
void
ListensSynchronizer::enqueListenNow(const Listen& listen)
{
enqueListen(listen, {});
}
void
@@ -313,15 +329,22 @@ namespace Scrobbling::ListenBrainz
if (timePoint.isValid())
{
const TimedListen timedListen {listen, timePoint};
// We want the listen to be sent again later in case of failure, so we just save it as pending send
const Database::ListenId listenId {saveListen(TimedListen {listen, timePoint})};
if (!listenId.isValid())
return;
saveListen(timedListen, Database::ScrobblingState::PendingAdd);
request.priority = Http::ClientRequestParameters::Priority::Normal;
request.onSuccessFunc = [=](std::string_view)
{
onListenSent(listenId);
_strand.dispatch([=]
{
if (saveListen(timedListen, Database::ScrobblingState::Synchronized))
{
UserContext& context {getUserContext(listen.userId)};
if (context.listenCount)
(*context.listenCount)++;
}
});
};
// on failure, this listen will be sent during the next sync
}
@@ -341,62 +364,93 @@ namespace Scrobbling::ListenBrainz
const std::optional<UUID> listenBrainzToken {Utils::getListenBrainzToken(_db.getTLSSession(), listen.userId)};
if (!listenBrainzToken)
{
LOG(DEBUG) << "No listenbrainz token found: skipping";
return;
}
request.message.addBodyText(bodyText);
request.message.addHeader("Authorization", "Token " + std::string {listenBrainzToken->getAsString()});
request.message.addHeader("Content-Type", "application/json");
_client.sendPOSTRequest(std::move(request));
}
Database::ListenId
ListensSynchronizer::saveListen(const TimedListen& listen)
bool
ListensSynchronizer::saveListen(const TimedListen& listen, Database::ScrobblingState scrobblingState)
{
using namespace Database;
Session& session {_db.getTLSSession()};
auto transaction {session.createUniqueTransaction()};
if (Database::Listen::find(session, listen.userId, listen.trackId, Database::Scrobbler::ListenBrainz, listen.listenedAt))
return {};
Database::Listen::pointer dbListen {Database::Listen::find(session, listen.userId, listen.trackId, Database::Scrobbler::ListenBrainz, listen.listenedAt)};
if (!dbListen)
{
const User::pointer user {User::find(session, listen.userId)};
if (!user)
return false;
const User::pointer user {User::find(session, listen.userId)};
if (!user)
return {};
const Track::pointer track {Track::find(session, listen.trackId)};
if (!track)
return false;
const Track::pointer track {Track::find(session, listen.trackId)};
if (!track)
return {};
dbListen = Database::Listen::create(session, user, track, Database::Scrobbler::ListenBrainz, listen.listenedAt);
dbListen.modify()->setScrobblingState(scrobblingState);
const auto dbListen {Database::Listen::create(session, user, track, Database::Scrobbler::ListenBrainz, listen.listenedAt)};
assert(dbListen->getScrobblingState() == Database::ScrobblingState::PendingAdd);
LOG(DEBUG) << "LISTEN CREATED for user " << user->getLoginName() << ", track '" << track->getName() << "' AT " << listen.listenedAt.toString();
return dbListen->getId();
return true;
}
if (dbListen->getScrobblingState() == scrobblingState)
return false;
dbListen.modify()->setScrobblingState(scrobblingState);
return true;
}
void
ListensSynchronizer::onListenSent(Database::ListenId listenId)
ListensSynchronizer::enquePendingListens()
{
_strand.dispatch([=]
std::vector<TimedListen> pendingListens;
{
Database::Session& session {_db.getTLSSession()};
auto transaction {session.createUniqueTransaction()};
if (Database::Listen::pointer listen {Database::Listen::find(session, listenId)})
{
listen.modify()->setScrobblingState(Database::ScrobblingState::Synchronized);
Database::Listen::FindParameters params;
params.setScrobbler(Database::Scrobbler::ListenBrainz)
.setScrobblingState(Database::ScrobblingState::PendingAdd)
.setRange(Database::Range {0, 100}); // don't flood too much?
UserContext& context {getUserContext(listen->getUser()->getId())};
if (context.listenCount)
(*context.listenCount)++;
const Database::RangeResults results {Database::Listen::find(session, params)};
pendingListens.reserve(results.results.size());
for (Database::ListenId listenId : results.results)
{
const Database::Listen::pointer listen {Database::Listen::find(session, listenId)};
TimedListen timedListen;
timedListen.listenedAt = listen->getDateTime();
timedListen.userId = listen->getUser()->getId();
timedListen.trackId = listen->getTrack()->getId();
pendingListens.push_back(std::move(timedListen));
}
});
}
LOG(DEBUG) << "Queing " << pendingListens.size() << " pending listen";
for (const TimedListen& pendingListen : pendingListens)
enqueListen(pendingListen);
}
ListensSynchronizer::UserContext&
ListensSynchronizer::getUserContext(Database::UserId userId)
{
assert(_strand.running_in_this_thread());
auto itContext {_userContexts.find(userId)};
if (itContext == std::cend(_userContexts))
{
@@ -408,24 +462,24 @@ namespace Scrobbling::ListenBrainz
}
bool
ListensSynchronizer::isFetching() const
ListensSynchronizer::isSyncing() const
{
return std::any_of(std::cbegin(_userContexts), std::cend(_userContexts), [](const auto& contextEntry)
{
const auto& [userId, context] {contextEntry};
return context.fetching;
return context.syncing;
});
}
void
ListensSynchronizer::scheduleGetListens(std::chrono::seconds fromNow)
ListensSynchronizer::scheduleSync(std::chrono::seconds fromNow)
{
if (_syncListensPeriod.count() == 0 || _maxSyncListenCount == 0)
return;
LOG(DEBUG) << "Scheduled sync in " << fromNow.count() << " seconds...";
_getListensTimer.expires_after(fromNow);
_getListensTimer.async_wait(boost::asio::bind_executor(_strand, [this] (const boost::system::error_code& ec)
_syncTimer.expires_after(fromNow);
_syncTimer.async_wait(boost::asio::bind_executor(_strand, [this] (const boost::system::error_code& ec)
{
if (ec == boost::asio::error::operation_aborted)
{
@@ -437,38 +491,37 @@ namespace Scrobbling::ListenBrainz
throw Exception {"GetListens timer failure: " + std::string {ec.message()} };
}
startGetListens();
startSync();
}));
}
void
ListensSynchronizer::startGetListens()
ListensSynchronizer::startSync()
{
LOG(DEBUG) << "GetListens started!";
LOG(DEBUG) << "Starting sync!";
assert(!isFetching());
assert(!isSyncing());
enquePendingListens();
Database::RangeResults<Database::UserId> userIds;
{
Database::Session& session {_db.getTLSSession()};
auto transaction {session.createSharedTransaction()};
userIds = Database::User::find(_db.getTLSSession(), Database::Range {});
userIds = Database::User::find(_db.getTLSSession(), Database::User::FindParameters{}.setScrobbler(Database::Scrobbler::ListenBrainz));
}
for (const Database::UserId userId : userIds.results)
{
if (Utils::getListenBrainzToken(_db.getTLSSession(), userId))
startGetListens(getUserContext(userId));
}
startSync(getUserContext(userId));
if (!isFetching())
scheduleGetListens(_syncListensPeriod);
if (!isSyncing())
scheduleSync(_syncListensPeriod);
}
void
ListensSynchronizer::startGetListens(UserContext& context)
ListensSynchronizer::startSync(UserContext& context)
{
context.fetching = true;
context.syncing = true;
context.listenBrainzUserName = "";
context.maxDateTime = {};
context.fetchedListenCount = 0;
@@ -479,15 +532,15 @@ namespace Scrobbling::ListenBrainz
}
void
ListensSynchronizer::onGetListensEnded(UserContext& context)
ListensSynchronizer::onSyncEnded(UserContext& context)
{
_strand.dispatch([this, &context]
{
LOG(DEBUG) << "Fetch done for user " << context.userId.getValue() << ", fetched: " << context.fetchedListenCount << ", matched: " << context.matchedListenCount << ", imported: " << context.importedListenCount;
context.fetching = false;
LOG_EX(context.importedListenCount > 0 ? Severity::INFO : Severity::DEBUG) << "Sync done for user '" << context.listenBrainzUserName << "', fetched: " << context.fetchedListenCount << ", matched: " << context.matchedListenCount << ", imported: " << context.importedListenCount;
context.syncing = false;
if (!isFetching())
scheduleGetListens(_syncListensPeriod);
if (!isSyncing())
scheduleSync(_syncListensPeriod);
});
}
@@ -499,7 +552,7 @@ namespace Scrobbling::ListenBrainz
const std::optional<UUID> listenBrainzToken {Utils::getListenBrainzToken(_db.getTLSSession(), context.userId)};
if (!listenBrainzToken)
{
onGetListensEnded(context);
onSyncEnded(context);
return;
}
@@ -512,14 +565,14 @@ namespace Scrobbling::ListenBrainz
context.listenBrainzUserName = parseValidateToken(msgBody);
if (context.listenBrainzUserName.empty())
{
onGetListensEnded(context);
onSyncEnded(context);
return;
}
enqueGetListenCount(context);
};
request.onFailureFunc = [this, &context]
{
onGetListensEnded(context);
onSyncEnded(context);
};
_client.sendGETRequest(std::move(request));
@@ -544,7 +597,7 @@ namespace Scrobbling::ListenBrainz
if (!needSync)
{
onGetListensEnded(context);
onSyncEnded(context);
return;
}
@@ -553,7 +606,7 @@ namespace Scrobbling::ListenBrainz
};
request.onFailureFunc = [this, &context]
{
onGetListensEnded(context);
onSyncEnded(context);
};
_client.sendGETRequest(std::move(request));
@@ -572,7 +625,7 @@ namespace Scrobbling::ListenBrainz
processGetListensResponse(msgBody, context);
if (context.fetchedListenCount >= _maxSyncListenCount || !context.maxDateTime.isValid())
{
onGetListensEnded(context);
onSyncEnded(context);
return;
}
@@ -580,7 +633,7 @@ namespace Scrobbling::ListenBrainz
};
request.onFailureFunc = [=, &context]
{
onGetListensEnded(context);
onSyncEnded(context);
};
_client.sendGETRequest(std::move(request));
@@ -597,27 +650,10 @@ namespace Scrobbling::ListenBrainz
context.matchedListenCount += parseResult.matchedListens.size();
context.maxDateTime = parseResult.oldestEntry;
if (parseResult.matchedListens.empty())
return;
auto transaction {session.createUniqueTransaction()};
Database::User::pointer user {Database::User::find(session, context.userId)};
if (!user)
return;
for (const TimedListen& listen : parseResult.matchedListens)
{
const Database::Track::pointer track {Database::Track::find(session, listen.trackId)};
if (!track)
continue;
if (Database::Listen::find(session, listen.userId, listen.trackId, Database::Scrobbler::ListenBrainz, listen.listenedAt))
continue;
Database::Listen::create(session, user, track, Database::Scrobbler::ListenBrainz, listen.listenedAt);
context.importedListenCount++;
if (saveListen(listen, Database::ScrobblingState::Synchronized))
context.importedListenCount++;
}
}
} // namespace Scrobbling::ListenBrainz
@@ -50,11 +50,14 @@ namespace Scrobbling::ListenBrainz
public:
ListensSynchronizer(boost::asio::io_context& ioContext, Database::Db& db, Http::IClient& client);
void enqueListen(const Listen& listen, const Wt::WDateTime& timePoint);
void enqueListen(const TimedListen& listen);
void enqueListenNow(const Listen& listen);
private:
Database::ListenId saveListen(const TimedListen& listen);
void onListenSent(Database::ListenId listenId);
void enqueListen(const Listen& listen, const Wt::WDateTime& timePoint);
bool saveListen(const TimedListen& listen, Database::ScrobblingState scrobblinState);
void enquePendingListens();
struct UserContext
{
@@ -66,10 +69,10 @@ namespace Scrobbling::ListenBrainz
UserContext& operator=(UserContext&&) = delete;
const Database::UserId userId;
bool fetching {};
bool syncing {};
std::optional<std::size_t> listenCount {};
// resetted at each fetch
// resetted at each sync
std::string listenBrainzUserName; // need to be resolved first
Wt::WDateTime maxDateTime;
std::size_t fetchedListenCount{};
@@ -78,20 +81,20 @@ namespace Scrobbling::ListenBrainz
};
UserContext& getUserContext(Database::UserId userId);
bool isFetching() const;
void scheduleGetListens(std::chrono::seconds fromNow);
void startGetListens();
void startGetListens(UserContext& context);
void onGetListensEnded(UserContext& context);
bool isSyncing() const;
void scheduleSync(std::chrono::seconds fromNow);
void startSync();
void startSync(UserContext& context);
void onSyncEnded(UserContext& context);
void enqueValidateToken(UserContext& context);
void enqueGetListenCount(UserContext& context);
void enqueGetListens(UserContext& context);
void processGetListensResponse(std::string_view body, UserContext& context);
void processGetListensResponse(std::string_view body, UserContext& context);
boost::asio::io_context& _ioContext;
boost::asio::io_context::strand _strand {_ioContext};
Database::Db& _db;
boost::asio::steady_timer _getListensTimer {_ioContext};
boost::asio::steady_timer _syncTimer {_ioContext};
Http::IClient& _client;
std::unordered_map<Database::UserId, UserContext> _userContexts;
@@ -19,14 +19,9 @@
#include "Utils.hpp"
#include <string_view>
#include "services/database/Session.hpp"
#include "services/database/TrackList.hpp"
#include "services/database/User.hpp"
static constexpr std::string_view historyTracklistName {"__scrobbler_listenbrainz_history__"};
namespace Scrobbling::ListenBrainz::Utils
{
std::optional<UUID>
@@ -38,9 +33,6 @@ namespace Scrobbling::ListenBrainz::Utils
if (!user)
return std::nullopt;
if (user->getScrobbler() != Database::Scrobbler::ListenBrainz)
return std::nullopt;
return user->getListenBrainzToken();
}
}
@@ -32,5 +32,5 @@ namespace Database
namespace Scrobbling::ListenBrainz::Utils
{
std::optional<UUID> getListenBrainzToken(Database::Session& session, Database::UserId userId);
std::optional<UUID> getListenBrainzToken(Database::Session& session, Database::UserId userId);
}