#include "poller.hpp"
#include <algorithm>
#include <chrono>
#include <optional>
#include <random>
#include <spdlog/spdlog.h>
#include "database.hpp"
#include "service_error.hpp"
namespace telegrammer
{
Poller::Poller(TelegramApi& client, DeliveryStore& store,
Dispatcher& dispatcher, RuntimeState& state, int64_t bot_id)
: client_(client), store_(store), dispatcher_(dispatcher), state_(state),
bot_id_(bot_id)
{}
void Poller::start()
{
worker_ = std::jthread([this](std::stop_token stop_token)
{
run(stop_token);
});
}
void Poller::setBotId(int64_t bot_id)
{
bot_id_ = bot_id;
}
void Poller::stop()
{
worker_.request_stop();
condition_.notify_all();
if(worker_.joinable())
{
worker_.join();
}
}
bool Poller::wait(std::stop_token stop_token,
std::chrono::milliseconds duration)
{
std::unique_lock lock(mutex_);
return condition_.wait_for(lock, stop_token, duration,
[] { return false; });
}
void Poller::run(std::stop_token stop_token)
{
int backoff_seconds = 1;
int consecutive_failures = 0;
auto first_failure = std::chrono::steady_clock::time_point{};
std::mt19937 random_generator(static_cast<unsigned>(
std::chrono::steady_clock::now().time_since_epoch().count()));
auto markFailure = [&](bool permanent)
{
auto current = std::chrono::steady_clock::now();
if(consecutive_failures == 0)
{
first_failure = current;
}
++consecutive_failures;
bool persistent = permanent || consecutive_failures >= 5 ||
current - first_failure >= std::chrono::minutes(5);
state_.degraded = persistent;
if(permanent)
{
state_.polling_ready = false;
}
};
auto markSuccess = [&]()
{
consecutive_failures = 0;
first_failure = {};
state_.polling_ready = true;
state_.degraded = false;
state_.last_success = nowSeconds();
};
auto waitBackoff = [&](std::optional<int> minimum_delay = std::nullopt)
{
int delay_seconds = backoff_seconds;
if(minimum_delay.has_value())
{
delay_seconds = std::max(delay_seconds, *minimum_delay);
}
std::uniform_int_distribution<int> jitter(
0, std::max(delay_seconds * 250, 1));
wait(stop_token,
std::chrono::seconds(delay_seconds) +
std::chrono::milliseconds(jitter(random_generator)));
backoff_seconds = std::min(backoff_seconds * 2, 60);
};
while(!stop_token.stop_requested())
{
try
{
auto offset = store_.offset(bot_id_);
if(!offset.has_value())
{
const ServiceError* error = asServiceError(offset.error());
bool permanent = error != nullptr && !error->retryable;
markFailure(permanent);
spdlog::error("Unable to read polling state: {}",
mw::errorMsg(offset.error()));
if(permanent)
{
break;
}
waitBackoff();
continue;
}
auto updates = client_.getUpdates(*offset, 30);
if(!updates.has_value())
{
const ServiceError* error = asServiceError(updates.error());
bool permanent = error != nullptr && !error->retryable;
markFailure(permanent);
if(error != nullptr)
{
spdlog::error("Telegram polling failed [{}]: {}",
error->code, error->msg);
if(permanent)
{
break;
}
}
else
{
spdlog::error("Telegram polling failed: {}",
mw::errorMsg(updates.error()));
}
waitBackoff(error == nullptr ? std::nullopt :
error->retry_after);
continue;
}
auto ingest_result = store_.ingest(bot_id_, (*updates)["result"]);
if(!ingest_result.has_value())
{
const ServiceError* error =
asServiceError(ingest_result.error());
bool permanent = error != nullptr &&
(!error->retryable ||
error->code == "QUEUE_FULL");
markFailure(permanent);
if(error != nullptr)
{
spdlog::error("Unable to ingest Telegram updates [{}]: {}",
error->code, error->msg);
}
else
{
spdlog::error("Unable to ingest Telegram updates: {}",
mw::errorMsg(ingest_result.error()));
}
waitBackoff();
continue;
}
markSuccess();
backoff_seconds = 1;
dispatcher_.notify();
// Add bounded jitter to avoid synchronized retries after an
// outage. A successful long poll normally does not wait here.
if((*updates)["result"].empty())
{
std::uniform_int_distribution<int> jitter(0, 250);
wait(stop_token,
std::chrono::seconds(0) +
std::chrono::milliseconds(jitter(random_generator)));
}
}
catch(const std::exception& error)
{
markFailure(false);
spdlog::error("Polling worker failed: {}", error.what());
waitBackoff();
}
catch(...)
{
markFailure(false);
spdlog::error("Polling worker failed with an unknown exception");
waitBackoff();
}
}
state_.polling_ready = false;
}
} // namespace telegrammer