feat(eventsub): properly unsubscribe once no more handles are interested (#5943)

This commit is contained in:
pajlada
2025-02-09 13:53:31 +01:00
committed by GitHub
parent 7dad35b7a7
commit 9527421bcd
6 changed files with 167 additions and 49 deletions
+1 -1
View File
@@ -22,7 +22,7 @@
- Bugfix: Fixed the reply button showing for inline whispers and announcements. (#5863) - Bugfix: Fixed the reply button showing for inline whispers and announcements. (#5863)
- Bugfix: Fixed suspicious user treatment update messages not being searchable. (#5865) - Bugfix: Fixed suspicious user treatment update messages not being searchable. (#5865)
- Bugfix: Ensure miniaudio backend exits even if it doesn't exit cleanly. (#5896) - Bugfix: Ensure miniaudio backend exits even if it doesn't exit cleanly. (#5896)
- Dev: Add initial experimental EventSub support. (#5837, #5895, #5897, #5904, #5910, #5903, #5915, #5916, #5930, #5935, #5932) - Dev: Add initial experimental EventSub support. (#5837, #5895, #5897, #5904, #5910, #5903, #5915, #5916, #5930, #5935, #5932, #5943)
- Dev: Remove unneeded platform specifier for toasts. (#5914) - Dev: Remove unneeded platform specifier for toasts. (#5914)
- Dev: Highlight checks now use non-capturing groups for the boundaries. (#5784) - Dev: Highlight checks now use non-capturing groups for the boundaries. (#5784)
- Dev: Removed unused PubSub whisper code. (#5898) - Dev: Removed unused PubSub whisper code. (#5898)
+5
View File
@@ -436,6 +436,11 @@ public:
failureCallback)), failureCallback)),
(override)); (override));
MOCK_METHOD(void, deleteEventSubSubscription,
(const QString &request, ResultCallback<> successCallback,
(FailureCallback<QString> failureCallback)),
(override));
MOCK_METHOD(void, update, (QString clientId, QString oauthToken), MOCK_METHOD(void, update, (QString clientId, QString oauthToken),
(override)); (override));
+40
View File
@@ -3334,6 +3334,46 @@ QDebug &operator<<(QDebug &dbg,
return dbg; return dbg;
} }
void Helix::deleteEventSubSubscription(const QString &subscriptionID,
ResultCallback<> successCallback,
FailureCallback<QString> failureCallback)
{
QUrlQuery query;
query.addQueryItem("id", subscriptionID);
this->makeDelete("eventsub/subscriptions", query)
.onSuccess([successCallback](const auto &result) {
if (result.status() != 204)
{
qCWarning(chatterinoTwitchEventSub)
<< "Success result for deleting eventsub subscription was "
<< result.formatError() << "but we expected it to be 204";
}
successCallback();
})
.onError([failureCallback](const NetworkResult &result) {
if (!result.status())
{
failureCallback(result.formatError());
return;
}
const auto obj = result.parseJson();
auto message = obj["message"].toString();
if (message.isEmpty())
{
failureCallback("Twitch internal server error");
}
else
{
failureCallback(message);
}
})
.execute();
}
NetworkRequest Helix::makeRequest(const QString &url, const QUrlQuery &urlQuery, NetworkRequest Helix::makeRequest(const QString &url, const QUrlQuery &urlQuery,
NetworkRequestType type) NetworkRequestType type)
{ {
+10
View File
@@ -1217,6 +1217,11 @@ public:
FailureCallback<HelixCreateEventSubSubscriptionError, QString> FailureCallback<HelixCreateEventSubSubscriptionError, QString>
failureCallback) = 0; failureCallback) = 0;
// https://dev.twitch.tv/docs/api/reference/#delete-eventsub-subscription
virtual void deleteEventSubSubscription(
const QString &subscriptionID, ResultCallback<> successCallback,
FailureCallback<QString> failureCallback) = 0;
virtual void update(QString clientId, QString oauthToken) = 0; virtual void update(QString clientId, QString oauthToken) = 0;
protected: protected:
@@ -1567,6 +1572,11 @@ public:
FailureCallback<HelixCreateEventSubSubscriptionError, QString> FailureCallback<HelixCreateEventSubSubscriptionError, QString>
failureCallback) final; failureCallback) final;
// https://dev.twitch.tv/docs/api/reference/#delete-eventsub-subscription
void deleteEventSubSubscription(
const QString &subscriptionID, ResultCallback<> successCallback,
FailureCallback<QString> failureCallback) final;
void update(QString clientId, QString oauthToken) final; void update(QString clientId, QString oauthToken) final;
static void initialize(); static void initialize();
+93 -38
View File
@@ -60,7 +60,6 @@ Controller::Controller()
Controller::~Controller() Controller::~Controller()
{ {
qCInfo(LOG) << "Controller dtor start"; qCInfo(LOG) << "Controller dtor start";
this->queuedSubscriptions.clear();
for (const auto &weakConnection : this->connections) for (const auto &weakConnection : this->connections)
{ {
@@ -73,6 +72,8 @@ Controller::~Controller()
connection->close(); connection->close();
} }
this->subscriptions.clear();
this->work.reset(); this->work.reset();
if (this->thread->joinable()) if (this->thread->joinable())
@@ -91,15 +92,35 @@ void Controller::removeRef(const SubscriptionRequest &request)
{ {
std::lock_guard lock(this->subscriptionsMutex); std::lock_guard lock(this->subscriptionsMutex);
assert(this->activeSubscriptions.contains(request)); assert(this->subscriptions.contains(request));
auto &xd = this->activeSubscriptions[request]; auto &subscription = this->subscriptions[request];
xd.refCount--; subscription.refCount--;
qCInfo(LOG) << "Remove ref for" << request << xd.refCount; qCInfo(LOG) << "Removed ref for" << request << subscription.refCount;
// todo use actual atomic things here to be smart
if (xd.refCount <= 0) if (subscription.refCount <= 0)
{ {
qCInfo(LOG) << "TODO: Unsubscribe from" << request; if (subscription.subscriptionID.isEmpty())
{
qCWarning(LOG) << "Refcount fell to zero for" << request
<< "but we had no subscription ID attached - was a "
"successful subscription never made?";
return;
}
qCInfo(LOG) << "Unsubscribing from" << request;
getHelix()->deleteEventSubSubscription(
subscription.subscriptionID,
[request] {
qCInfo(LOG) << "Successfully unsubscribed from" << request;
},
[request](const auto &errorMessage) {
qCInfo(LOG)
<< "An error occurred while attempting to unsubscribe from"
<< request << errorMessage;
});
subscription.subscriptionID.clear();
} }
} }
@@ -110,13 +131,13 @@ SubscriptionHandle Controller::subscribe(const SubscriptionRequest &request)
{ {
std::lock_guard lock(this->subscriptionsMutex); std::lock_guard lock(this->subscriptionsMutex);
auto &xd = this->activeSubscriptions[request]; auto &subscription = this->subscriptions[request];
if (xd.refCount == 0) if (subscription.refCount == 0)
{ {
needToSubscribe = true; needToSubscribe = true;
} }
xd.refCount++; subscription.refCount++;
qCInfo(LOG) << "Add ref for" << request << xd.refCount; qCInfo(LOG) << "Added ref for" << request << subscription.refCount;
} }
auto handle = std::make_unique<RawSubscriptionHandle>(request); auto handle = std::make_unique<RawSubscriptionHandle>(request);
@@ -131,24 +152,29 @@ SubscriptionHandle Controller::subscribe(const SubscriptionRequest &request)
return handle; return handle;
} }
void Controller::subscribe(const SubscriptionRequest &request, bool isQueued) void Controller::subscribe(const SubscriptionRequest &request, bool isRetry)
{ {
qCInfo(LOG) << "Subscribe request for" << request.subscriptionType; qCInfo(LOG) << "Subscribe request for" << request.subscriptionType;
boost::asio::post(this->ioContext, [this, request, isQueued] { boost::asio::post(this->ioContext, [this, request, isRetry] {
// 1. Flush dead connections (maybe this should not be done here) // 1. Flush dead connections (maybe this should not be done here)
// TODO: implement // TODO: implement
if (isQueued)
{ {
qCInfo(LOG) << "Removing subscription from queued list"; std::lock_guard lock(this->subscriptionsMutex);
this->queuedSubscriptions.erase(request); auto &subscription = this->subscriptions[request];
} if (isRetry)
{
qCInfo(LOG) << "Removing subscription from queued list";
if (this->queuedSubscriptions.contains(request)) subscription.retryTimer.reset();
{ }
qCWarning(LOG) << "We already have a queued subscription for this, " else if (subscription.retryTimer)
"let's chill :)"; {
return; qCWarning(LOG)
<< "We already have a queued subscription for this, "
"let's chill :)";
return;
}
} }
uint32_t openButNotReadyConnections = 0; uint32_t openButNotReadyConnections = 0;
@@ -188,7 +214,8 @@ void Controller::subscribe(const SubscriptionRequest &request, bool isQueued)
[this, request, connection, [this, request, connection,
weakConnection{weakConnection}](const auto &res) { weakConnection{weakConnection}](const auto &res) {
qCInfo(LOG) << "success" << res; qCInfo(LOG) << "success" << res;
this->markRequestSubscribed(request, weakConnection); this->markRequestSubscribed(request, weakConnection,
res.subscriptionID);
}, },
[this, request](const auto &error, const auto &errorString) { [this, request](const auto &error, const auto &errorString) {
using Error = HelixCreateEventSubSubscriptionError; using Error = HelixCreateEventSubSubscriptionError;
@@ -204,6 +231,10 @@ void Controller::subscribe(const SubscriptionRequest &request, bool isQueued)
case Error::Forbidden: case Error::Forbidden:
qCWarning(LOG) << "Forbidden" << errorString; qCWarning(LOG) << "Forbidden" << errorString;
boost::asio::post(this->ioContext, [this, request] {
this->retrySubscription(
request, boost::posix_time::seconds(2), 5);
});
break; break;
case Error::Conflict: case Error::Conflict:
@@ -221,8 +252,8 @@ void Controller::subscribe(const SubscriptionRequest &request, bool isQueued)
<< "Unhandled error, retrying subscription" << "Unhandled error, retrying subscription"
<< errorString; << errorString;
boost::asio::post(this->ioContext, [this, request] { boost::asio::post(this->ioContext, [this, request] {
this->queueSubscription( this->retrySubscription(
request, boost::posix_time::seconds(2)); request, boost::posix_time::seconds(2), 5);
}); });
break; break;
} }
@@ -234,18 +265,20 @@ void Controller::subscribe(const SubscriptionRequest &request, bool isQueued)
{ {
// No connection was available to handle this subscription request, create a new connection // No connection was available to handle this subscription request, create a new connection
this->createConnection(); this->createConnection();
this->queueSubscription(request, boost::posix_time::millisec(500)); this->retrySubscription(request, boost::posix_time::millisec(500),
10);
} }
else else
{ {
if (openButNotReadyConnections > 1) if (openButNotReadyConnections > 1)
{ {
qCWarning(LOG) << "We have" << openButNotReadyConnections qCWarning(LOG) << "We have" << openButNotReadyConnections
<< "open but not ready connections, hmmm"; << "open but no ready connections, hmmm";
} }
// At least one connection is open, but it has not gotten the welcome message yet // At least one connection is open, but it has not gotten the welcome message yet
this->queueSubscription(request, boost::posix_time::millisec(250)); this->retrySubscription(request, boost::posix_time::millisec(250),
10);
} }
}); });
} }
@@ -288,15 +321,30 @@ void Controller::registerConnection(std::weak_ptr<lib::Session> &&connection)
this->connections.emplace_back(std::move(connection)); this->connections.emplace_back(std::move(connection));
} }
void Controller::queueSubscription(const SubscriptionRequest &request, void Controller::retrySubscription(const SubscriptionRequest &request,
boost::posix_time::time_duration delay) boost::posix_time::time_duration delay,
int32_t maxAttempts)
{ {
this->threadGuard->guard(); std::lock_guard lock(this->subscriptionsMutex);
auto resubTimer = auto &connection = this->subscriptions[request];
if (connection.retryAttempts <= 0)
{
connection.retryAttempts = maxAttempts;
}
else if (--connection.retryAttempts == 0)
{
qCWarning(LOG) << "Reached max amount of retries for" << request;
return;
}
qCInfo(LOG) << "Retrying subscription" << request << " - attempt"
<< connection.retryAttempts;
auto retryTimer =
std::make_unique<boost::asio::deadline_timer>(this->ioContext); std::make_unique<boost::asio::deadline_timer>(this->ioContext);
resubTimer->expires_from_now(delay); retryTimer->expires_from_now(delay);
resubTimer->async_wait([this, request](const auto &ec) { retryTimer->async_wait([this, request](const auto &ec) {
if (!ec) if (!ec)
{ {
// The timer passed naturally // The timer passed naturally
@@ -304,15 +352,22 @@ void Controller::queueSubscription(const SubscriptionRequest &request,
} }
}); });
this->queuedSubscriptions.emplace(request, std::move(resubTimer)); assert(connection.retryTimer == nullptr &&
"Timer should not already be set");
connection.retryTimer = std::move(retryTimer);
} }
void Controller::markRequestSubscribed(const SubscriptionRequest &request, void Controller::markRequestSubscribed(const SubscriptionRequest &request,
std::weak_ptr<lib::Session> connection) std::weak_ptr<lib::Session> connection,
const QString &subscriptionID)
{ {
std::lock_guard lock(this->subscriptionsMutex); std::lock_guard lock(this->subscriptionsMutex);
this->activeSubscriptions[request].connection = std::move(connection); auto &subscription = this->subscriptions[request];
subscription.connection = std::move(connection);
subscription.subscriptionID = subscriptionID;
} }
} // namespace chatterino::eventsub } // namespace chatterino::eventsub
+18 -10
View File
@@ -22,6 +22,9 @@ class IController
public: public:
virtual ~IController() = default; virtual ~IController() = default;
/// Removes one reference for the given subscription request
///
/// Should realistically only be called in the dtor of SubscriptionHandle
virtual void removeRef(const SubscriptionRequest &request) = 0; virtual void removeRef(const SubscriptionRequest &request) = 0;
/// Subscribe will make a request to each open connection and ask them to /// Subscribe will make a request to each open connection and ask them to
@@ -47,16 +50,18 @@ public:
const SubscriptionRequest &request) override; const SubscriptionRequest &request) override;
private: private:
void subscribe(const SubscriptionRequest &request, bool isQueued); void subscribe(const SubscriptionRequest &request, bool isRetry);
void createConnection(); void createConnection();
void registerConnection(std::weak_ptr<lib::Session> &&connection); void registerConnection(std::weak_ptr<lib::Session> &&connection);
void queueSubscription(const SubscriptionRequest &request, void retrySubscription(const SubscriptionRequest &request,
boost::posix_time::time_duration delay); boost::posix_time::time_duration delay,
int32_t maxAttempts);
void markRequestSubscribed(const SubscriptionRequest &request, void markRequestSubscribed(const SubscriptionRequest &request,
std::weak_ptr<lib::Session> connection); std::weak_ptr<lib::Session> connection,
const QString &subscriptionID);
const std::string userAgent; const std::string userAgent;
@@ -72,17 +77,20 @@ private:
std::vector<std::weak_ptr<lib::Session>> connections; std::vector<std::weak_ptr<lib::Session>> connections;
struct XD { struct Subscription {
int32_t refCount = 0; int32_t refCount = 0;
std::weak_ptr<lib::Session> connection; std::weak_ptr<lib::Session> connection;
/// The ID of the subscription the Twitch Helix API has given us
QString subscriptionID;
/// The timer, if any, for retrying the subscription creation
std::unique_ptr<boost::asio::deadline_timer> retryTimer;
int32_t retryAttempts = 0;
}; };
std::mutex subscriptionsMutex; std::mutex subscriptionsMutex;
std::unordered_map<SubscriptionRequest, XD> activeSubscriptions; std::unordered_map<SubscriptionRequest, Subscription> subscriptions;
std::unordered_map<SubscriptionRequest,
std::unique_ptr<boost::asio::deadline_timer>>
queuedSubscriptions;
}; };
class DummyController : public IController class DummyController : public IController