Changed http client to be instanciated by services
This commit is contained in:
@@ -30,6 +30,7 @@
|
||||
#include "services/database/StarredTrack.hpp"
|
||||
#include "services/database/Track.hpp"
|
||||
#include "services/database/User.hpp"
|
||||
#include "utils/Logger.hpp"
|
||||
|
||||
#include "internal/InternalScrobbler.hpp"
|
||||
#include "listenbrainz/ListenBrainzScrobbler.hpp"
|
||||
@@ -47,8 +48,15 @@ namespace Scrobbling
|
||||
ScrobblingService::ScrobblingService(boost::asio::io_context& ioContext, Db& db)
|
||||
: _db {db}
|
||||
{
|
||||
LMS_LOG(SCROBBLING, INFO) << "Starting service...";
|
||||
_scrobblers.emplace(Scrobbler::Internal, std::make_unique<InternalScrobbler>(_db));
|
||||
_scrobblers.emplace(Scrobbler::ListenBrainz, std::make_unique<ListenBrainz::Scrobbler>(ioContext, _db));
|
||||
LMS_LOG(SCROBBLING, INFO) << "Service started!";
|
||||
}
|
||||
|
||||
ScrobblingService::~ScrobblingService()
|
||||
{
|
||||
LMS_LOG(SCROBBLING, INFO) << "Service stopped!";
|
||||
}
|
||||
|
||||
void
|
||||
|
||||
@@ -32,6 +32,7 @@ namespace Scrobbling
|
||||
{
|
||||
public:
|
||||
ScrobblingService(boost::asio::io_context& ioContext, Database::Db& db);
|
||||
~ScrobblingService();
|
||||
|
||||
private:
|
||||
void listenStarted(const Listen& listen) override;
|
||||
|
||||
@@ -144,7 +144,8 @@ namespace Scrobbling::ListenBrainz
|
||||
: _ioContext {ioContext}
|
||||
, _db {db}
|
||||
, _baseAPIUrl {Service<IConfig>::get()->getString("listenbrainz-api-base-url", "https://api.listenbrainz.org")}
|
||||
, _listensSynchronizer {_ioContext, db, _baseAPIUrl}
|
||||
, _client {Http::createClient(_ioContext, _baseAPIUrl)}
|
||||
, _listensSynchronizer {_ioContext, db, *_client}
|
||||
{
|
||||
LOG(INFO) << "Starting ListenBrainz scrobbler... API endpoint = '" << _baseAPIUrl;
|
||||
}
|
||||
@@ -183,7 +184,7 @@ namespace Scrobbling::ListenBrainz
|
||||
Scrobbler::enqueListen(const Listen& listen, const Wt::WDateTime& timePoint)
|
||||
{
|
||||
Http::ClientPOSTRequestParameters request;
|
||||
request.url = _baseAPIUrl + "/1/submit-listens";
|
||||
request.relativeUrl = "/1/submit-listens";
|
||||
|
||||
if (timePoint.isValid())
|
||||
{
|
||||
@@ -213,7 +214,7 @@ namespace Scrobbling::ListenBrainz
|
||||
request.message.addBodyText(bodyText);
|
||||
request.message.addHeader("Authorization", "Token " + std::string {listenBrainzToken->getAsString()});
|
||||
request.message.addHeader("Content-Type", "application/json");
|
||||
Service<Http::IClient>::get()->sendPOSTRequest(std::move(request));
|
||||
_client->sendPOSTRequest(std::move(request));
|
||||
}
|
||||
} // namespace Scrobbling::ListenBrainz
|
||||
|
||||
|
||||
@@ -53,10 +53,11 @@ namespace Scrobbling::ListenBrainz
|
||||
// Submit listens
|
||||
void enqueListen(const Listen& listen, const Wt::WDateTime& timePoint);
|
||||
|
||||
boost::asio::io_context& _ioContext;
|
||||
Database::Db& _db;
|
||||
std::string _baseAPIUrl;
|
||||
ListensSynchronizer _listensSynchronizer;
|
||||
boost::asio::io_context& _ioContext;
|
||||
Database::Db& _db;
|
||||
std::string _baseAPIUrl;
|
||||
std::unique_ptr<Http::IClient> _client;
|
||||
ListensSynchronizer _listensSynchronizer;
|
||||
};
|
||||
} // Scrobbling::ListenBrainz
|
||||
|
||||
|
||||
@@ -213,10 +213,10 @@ namespace
|
||||
|
||||
namespace Scrobbling::ListenBrainz
|
||||
{
|
||||
ListensSynchronizer::ListensSynchronizer(boost::asio::io_context& ioContext, Database::Db& db, std::string_view baseAPIUrl)
|
||||
ListensSynchronizer::ListensSynchronizer(boost::asio::io_context& ioContext, Database::Db& db, Http::IClient& client)
|
||||
: _ioContext {ioContext}
|
||||
, _db {db}
|
||||
, _baseAPIUrl {baseAPIUrl}
|
||||
, _client {client}
|
||||
, _maxSyncListenCount {Service<IConfig>::get()->getULong("listenbrainz-max-sync-listen-count", 1000)}
|
||||
, _syncListensPeriod {Service<IConfig>::get()->getULong("listenbrainz-sync-listens-period-hours", 1)}
|
||||
{
|
||||
@@ -364,7 +364,7 @@ namespace Scrobbling::ListenBrainz
|
||||
|
||||
Http::ClientGETRequestParameters request;
|
||||
request.priority = Http::ClientRequestParameters::Priority::Low;
|
||||
request.url = _baseAPIUrl + "/1/validate-token";
|
||||
request.relativeUrl = "/1/validate-token";
|
||||
request.headers = { {"Authorization", "Token " + std::string {listenBrainzToken->getAsString()}} };
|
||||
request.onSuccessFunc = [this, &context] (std::string_view msgBody)
|
||||
{
|
||||
@@ -381,7 +381,7 @@ namespace Scrobbling::ListenBrainz
|
||||
onGetListensEnded(context);
|
||||
};
|
||||
|
||||
Service<Http::IClient>::get()->sendGETRequest(std::move(request));
|
||||
_client.sendGETRequest(std::move(request));
|
||||
}
|
||||
|
||||
void
|
||||
@@ -390,7 +390,7 @@ namespace Scrobbling::ListenBrainz
|
||||
assert(!context.listenBrainzUserName.empty());
|
||||
|
||||
Http::ClientGETRequestParameters request;
|
||||
request.url = _baseAPIUrl + "/1/user/" + std::string {context.listenBrainzUserName} + "/listen-count";
|
||||
request.relativeUrl = "/1/user/" + std::string {context.listenBrainzUserName} + "/listen-count";
|
||||
request.priority = Http::ClientRequestParameters::Priority::Low;
|
||||
request.onSuccessFunc = [=, &context] (std::string_view msgBody)
|
||||
{
|
||||
@@ -415,7 +415,7 @@ namespace Scrobbling::ListenBrainz
|
||||
onGetListensEnded(context);
|
||||
};
|
||||
|
||||
Service<Http::IClient>::get()->sendGETRequest(std::move(request));
|
||||
_client.sendGETRequest(std::move(request));
|
||||
}
|
||||
|
||||
void
|
||||
@@ -424,7 +424,7 @@ namespace Scrobbling::ListenBrainz
|
||||
assert(!context.listenBrainzUserName.empty());
|
||||
|
||||
Http::ClientGETRequestParameters request;
|
||||
request.url = _baseAPIUrl + "/1/user/" + context.listenBrainzUserName + "/listens?max_ts=" + std::to_string(context.maxDateTime.toTime_t());
|
||||
request.relativeUrl = "/1/user/" + context.listenBrainzUserName + "/listens?max_ts=" + std::to_string(context.maxDateTime.toTime_t());
|
||||
request.priority = Http::ClientRequestParameters::Priority::Low;
|
||||
request.onSuccessFunc = [=, &context] (std::string_view msgBody)
|
||||
{
|
||||
@@ -442,7 +442,7 @@ namespace Scrobbling::ListenBrainz
|
||||
onGetListensEnded(context);
|
||||
};
|
||||
|
||||
Service<Http::IClient>::get()->sendGETRequest(std::move(request));
|
||||
_client.sendGETRequest(std::move(request));
|
||||
}
|
||||
|
||||
void
|
||||
|
||||
@@ -36,12 +36,17 @@ namespace Database
|
||||
class User;
|
||||
}
|
||||
|
||||
namespace Http
|
||||
{
|
||||
class IClient;
|
||||
}
|
||||
|
||||
namespace Scrobbling::ListenBrainz
|
||||
{
|
||||
class ListensSynchronizer
|
||||
{
|
||||
public:
|
||||
ListensSynchronizer(boost::asio::io_context& ioContext, Database::Db& db, std::string_view baseAPIUrl);
|
||||
ListensSynchronizer(boost::asio::io_context& ioContext, Database::Db& db, Http::IClient& client);
|
||||
|
||||
void saveListen(const TimedListen& listen);
|
||||
|
||||
@@ -81,8 +86,8 @@ namespace Scrobbling::ListenBrainz
|
||||
boost::asio::io_context& _ioContext;
|
||||
boost::asio::io_context::strand _strand {_ioContext};
|
||||
Database::Db& _db;
|
||||
std::string _baseAPIUrl;
|
||||
boost::asio::steady_timer _getListensTimer {_ioContext};
|
||||
Http::IClient& _client;
|
||||
|
||||
std::unordered_map<Database::UserId, UserContext> _userContexts;
|
||||
|
||||
|
||||
@@ -27,7 +27,7 @@ IOContextRunner::IOContextRunner(boost::asio::io_service& ioService, std::size_t
|
||||
: _ioService {ioService}
|
||||
, _work {ioService}
|
||||
{
|
||||
LMS_LOG(UTILS, INFO) << "Starting IO Context with " << threadCount << " threads...";
|
||||
LMS_LOG(UTILS, INFO) << "Starting IO context with " << threadCount << " threads...";
|
||||
for (std::size_t i {}; i < threadCount; ++i)
|
||||
{
|
||||
_threads.emplace_back([&]
|
||||
@@ -48,10 +48,10 @@ IOContextRunner::IOContextRunner(boost::asio::io_service& ioService, std::size_t
|
||||
void
|
||||
IOContextRunner::stop()
|
||||
{
|
||||
LMS_LOG(UTILS, INFO) << "Stopping IO Context";
|
||||
LMS_LOG(UTILS, DEBUG) << "Stopping IO context...";
|
||||
_work.reset();
|
||||
_ioService.stop();
|
||||
LMS_LOG(UTILS, INFO) << "Stopped IO Context";
|
||||
LMS_LOG(UTILS, DEBUG) << "IO context stopped!";
|
||||
}
|
||||
|
||||
IOContextRunner::~IOContextRunner()
|
||||
|
||||
@@ -23,51 +23,21 @@
|
||||
namespace Http
|
||||
{
|
||||
std::unique_ptr<IClient>
|
||||
createClient(boost::asio::io_context& ioContext)
|
||||
createClient(boost::asio::io_context& ioContext, std::string_view baseUrl)
|
||||
{
|
||||
return std::make_unique<Client>(ioContext);
|
||||
return std::make_unique<Client>(ioContext, baseUrl);
|
||||
}
|
||||
|
||||
void
|
||||
Client::sendGETRequest(ClientGETRequestParameters&& GETParams)
|
||||
{
|
||||
SendQueue& sendQueue {getOrCreateSendQueue(GETParams.url)};
|
||||
sendQueue.sendRequest(std::make_unique<ClientRequest>(std::move(GETParams)));
|
||||
_sendQueue.sendRequest(std::make_unique<ClientRequest>(std::move(GETParams)));
|
||||
}
|
||||
|
||||
void
|
||||
Client::sendPOSTRequest(ClientPOSTRequestParameters&& POSTParams)
|
||||
{
|
||||
SendQueue& sendQueue {getOrCreateSendQueue(POSTParams.url)};
|
||||
sendQueue.sendRequest(std::make_unique<ClientRequest>(std::move(POSTParams)));
|
||||
_sendQueue.sendRequest(std::make_unique<ClientRequest>(std::move(POSTParams)));
|
||||
}
|
||||
|
||||
SendQueue&
|
||||
Client::getOrCreateSendQueue(const std::string& url)
|
||||
{
|
||||
Wt::Http::Client::URL parsedURL;
|
||||
if (!Wt::Http::Client::parseUrl(url, parsedURL))
|
||||
throw LmsException {"Cannot parse URL '" + url + "'"};
|
||||
|
||||
{
|
||||
std::shared_lock lock {_sendQueuesMutex};
|
||||
|
||||
if (auto it = _sendQueues.find(parsedURL.host); it != std::cend(_sendQueues))
|
||||
return it->second;
|
||||
}
|
||||
|
||||
{
|
||||
std::unique_lock lock {_sendQueuesMutex};
|
||||
|
||||
if (auto it = _sendQueues.find(parsedURL.host); it != std::cend(_sendQueues))
|
||||
return it->second;
|
||||
|
||||
auto [it, inserted] {_sendQueues.emplace(parsedURL.host, _ioContext)};
|
||||
assert(inserted);
|
||||
|
||||
return it->second;
|
||||
}
|
||||
}
|
||||
|
||||
} // namespace Http
|
||||
|
||||
|
||||
@@ -31,17 +31,17 @@ namespace Http
|
||||
class Client final : public IClient
|
||||
{
|
||||
public:
|
||||
Client(boost::asio::io_context& ioContext) : _ioContext {ioContext} {}
|
||||
Client(boost::asio::io_context& ioContext, std::string_view baseUrl)
|
||||
: _ioContext {ioContext}
|
||||
, _sendQueue {ioContext, baseUrl}
|
||||
{}
|
||||
|
||||
private:
|
||||
void sendGETRequest(ClientGETRequestParameters&& request) override;
|
||||
void sendPOSTRequest(ClientPOSTRequestParameters&& request) override;
|
||||
|
||||
SendQueue& getOrCreateSendQueue(const std::string& host);
|
||||
|
||||
boost::asio::io_context& _ioContext;
|
||||
std::shared_mutex _sendQueuesMutex;
|
||||
std::unordered_map<std::string, SendQueue> _sendQueues;
|
||||
boost::asio::io_context& _ioContext;
|
||||
SendQueue _sendQueue;
|
||||
};
|
||||
} // namespace Http
|
||||
|
||||
|
||||
@@ -60,8 +60,9 @@ namespace
|
||||
|
||||
namespace Http
|
||||
{
|
||||
SendQueue::SendQueue(boost::asio::io_context& ioContext)
|
||||
SendQueue::SendQueue(boost::asio::io_context& ioContext, std::string_view baseUrl)
|
||||
: _ioContext {ioContext}
|
||||
, _baseUrl {baseUrl}
|
||||
{
|
||||
_client.done().connect([this](Wt::AsioWrapper::error_code ec, const Wt::Http::Message& msg)
|
||||
{
|
||||
@@ -116,17 +117,18 @@ namespace Http
|
||||
bool
|
||||
SendQueue::sendRequest(const ClientRequest& request)
|
||||
{
|
||||
LOG(DEBUG) << "Sending request to url '" << request.getParameters().url << "'";
|
||||
std::string url {_baseUrl + request.getParameters().relativeUrl};
|
||||
LOG(DEBUG) << "Sending request to url '" << url << "'";
|
||||
|
||||
bool res {};
|
||||
switch (request.getType())
|
||||
{
|
||||
case ClientRequest::Type::GET:
|
||||
res = _client.get(request.getParameters().url, request.getGETParameters().headers);
|
||||
res = _client.get(url, request.getGETParameters().headers);
|
||||
break;
|
||||
|
||||
case ClientRequest::Type::POST:
|
||||
res = _client.post(request.getParameters().url, request.getPOSTParameters().message);
|
||||
res = _client.post(url, request.getPOSTParameters().message);
|
||||
break;
|
||||
}
|
||||
|
||||
|
||||
@@ -35,7 +35,7 @@ namespace Http
|
||||
class SendQueue
|
||||
{
|
||||
public:
|
||||
SendQueue(boost::asio::io_context& ioContext);
|
||||
SendQueue(boost::asio::io_context& ioContext, std::string_view baseUrl);
|
||||
~SendQueue();
|
||||
|
||||
SendQueue(const SendQueue&) = delete;
|
||||
@@ -62,6 +62,7 @@ namespace Http
|
||||
boost::asio::io_context& _ioContext;
|
||||
boost::asio::io_context::strand _strand {_ioContext};
|
||||
boost::asio::steady_timer _throttleTimer {_ioContext};
|
||||
std::string _baseUrl;
|
||||
|
||||
enum class State
|
||||
{
|
||||
|
||||
@@ -37,7 +37,7 @@ namespace Http
|
||||
};
|
||||
|
||||
Priority priority {Priority::Normal};
|
||||
std::string url;
|
||||
std::string relativeUrl; // relative to baseUrl used by the client
|
||||
|
||||
using OnSuccessFunc = std::function<void(std::string_view msgBody)>;
|
||||
OnSuccessFunc onSuccessFunc;
|
||||
|
||||
@@ -35,6 +35,6 @@ namespace Http
|
||||
virtual void sendPOSTRequest(ClientPOSTRequestParameters&& request) = 0;
|
||||
};
|
||||
|
||||
std::unique_ptr<IClient> createClient(boost::asio::io_context& ioContext);
|
||||
std::unique_ptr<IClient> createClient(boost::asio::io_context& ioContext, std::string_view baseUrl);
|
||||
} // namespace Http
|
||||
|
||||
|
||||
@@ -37,7 +37,6 @@
|
||||
#include "subsonic/SubsonicResource.hpp"
|
||||
#include "ui/LmsApplication.hpp"
|
||||
#include "ui/LmsApplicationManager.hpp"
|
||||
#include "utils/http/IClient.hpp"
|
||||
#include "utils/IChildProcessManager.hpp"
|
||||
#include "utils/IConfig.hpp"
|
||||
#include "utils/IOContextRunner.hpp"
|
||||
@@ -255,7 +254,6 @@ int main(int argc, char* argv[])
|
||||
else
|
||||
throw LmsException {"Bad value '" + authenticationBackend + "' for 'authentication-backend'"};
|
||||
|
||||
Service<Http::IClient> httpClient {Http::createClient(ioContext)};
|
||||
Service<Cover::ICoverService> coverService {Cover::createCoverService(database, argv[0], server.appRoot() + "/images/unknown-cover.jpg")};
|
||||
Service<Recommendation::IRecommendationService> recommendationService {Recommendation::createRecommendationService(database)};
|
||||
Service<Scanner::IScannerService> scannerService {Scanner::createScannerService(database, *recommendationService)};
|
||||
|
||||
Reference in New Issue
Block a user