#include "dispatcher.hpp"
#include <chrono>
#include <cctype>
#include <ctime>
#include <format>
#include <optional>
#include <string>
#include <curl/curl.h>
#include <mw/http_client.hpp>
#include <spdlog/spdlog.h>
#include "service_error.hpp"
namespace telegrammer
{
namespace
{
std::optional<int> retryAfterDate(const std::string& value)
{
time_t retry_time = curl_getdate(value.c_str(), nullptr);
time_t current_time = std::time(nullptr);
if(retry_time < current_time || retry_time - current_time > 86400)
{
return std::nullopt;
}
return static_cast<int>(retry_time - current_time);
}
std::optional<int> retryAfter(const mw::HTTPResponse& response)
{
auto iterator = response.header.find("Retry-After");
if(iterator == response.header.end())
{
iterator = response.header.find("retry-after");
}
if(iterator == response.header.end())
{
return std::nullopt;
}
try
{
std::size_t position = 0;
int value = std::stoi(iterator->second, &position);
if(position != iterator->second.size() || value < 0 || value > 86400)
{
return retryAfterDate(iterator->second);
}
return value;
}
catch(const std::exception&)
{
return retryAfterDate(iterator->second);
}
}
std::string deliveryId(const DeliveryJob& job)
{
return std::format("{}-{}-{}", job.bot_id, job.subscription_id,
job.update_id);
}
} // namespace
Dispatcher::Dispatcher(DeliveryStore& store, std::size_t worker_count)
: store_(store), worker_count_(worker_count)
{}
bool Dispatcher::shouldPurge()
{
std::lock_guard lock(purge_mutex_);
auto now = std::chrono::steady_clock::now();
if(now < next_purge_)
{
return false;
}
next_purge_ = now + std::chrono::hours(1);
return true;
}
void Dispatcher::start()
{
for(std::size_t i = 0; i < worker_count_; ++i)
{
workers_.emplace_back([this](std::stop_token stop_token)
{
run(stop_token);
});
}
}
void Dispatcher::stop()
{
for(std::jthread& worker: workers_)
{
worker.request_stop();
}
condition_.notify_all();
workers_.clear();
}
void Dispatcher::notify()
{
condition_.notify_all();
}
void Dispatcher::run(std::stop_token stop_token)
{
while(!stop_token.stop_requested())
{
if(shouldPurge())
{
auto purged = store_.purgeExpiredDead();
if(!purged.has_value())
{
spdlog::error("Unable to purge dead callback jobs: {}",
mw::errorMsg(purged.error()));
}
else if(*purged > 0)
{
spdlog::warn("Purged {} expired dead callback jobs", *purged);
}
}
std::optional<DeliveryJob> active_job;
try
{
auto result = store_.claimNext();
if(!result.has_value())
{
spdlog::error("Unable to claim callback job: {}",
mw::errorMsg(result.error()));
std::unique_lock lock(mutex_);
condition_.wait_for(lock, stop_token,
std::chrono::seconds(1),
[] { return false; });
continue;
}
if(!result->has_value())
{
std::unique_lock lock(mutex_);
condition_.wait_for(lock, stop_token,
std::chrono::seconds(1),
[] { return false; });
continue;
}
DeliveryJob job = std::move(result->value());
active_job = job;
mw::HTTPSession session;
auto connection = session.connectionTimeout(
std::chrono::seconds(3));
auto transfer = session.transferTimeout(std::chrono::seconds(10));
auto size = session.maxSize(64 * 1024);
auto protocols = session.allowedProtocols("http,https");
session.followRedirects(false);
if(!connection.has_value() || !transfer.has_value() ||
!size.has_value() || !protocols.has_value())
{
auto failure = store_.fail(
job.id, "Unable to configure callback HTTP client", true,
std::nullopt);
if(!failure.has_value())
{
spdlog::error("Unable to record callback failure {}: {}",
job.id, mw::errorMsg(failure.error()));
}
active_job.reset();
continue;
}
mw::HTTPRequest request(job.callback_url);
request.setContentType("application/json");
request.addHeader("X-Telegrammer-Delivery-Id", deliveryId(job));
request.setPayload(job.payload);
auto response = session.post(request);
if(response.has_value() && (*response)->status >= 200 &&
(*response)->status < 300)
{
auto complete = store_.complete(job.id);
if(!complete.has_value())
{
spdlog::error("Unable to complete callback job {}: {}",
job.id, mw::errorMsg(complete.error()));
auto failure = store_.fail(
job.id, "Unable to record callback completion", true,
std::nullopt);
if(!failure.has_value())
{
spdlog::error("Unable to recover callback job {}: {}",
job.id, mw::errorMsg(failure.error()));
}
}
active_job.reset();
continue;
}
bool should_retry = true;
std::optional<int> retry_after;
std::string failure_message = "Callback request failed";
if(response.has_value())
{
int status = (*response)->status;
should_retry = status == 408 || status == 429 || status >= 500;
retry_after = retryAfter(**response);
failure_message = std::format("Callback returned HTTP {}",
status);
}
else
{
failure_message = "Callback transport failed";
}
auto failure = store_.fail(job.id, failure_message, should_retry,
retry_after);
if(!failure.has_value())
{
spdlog::error("Unable to record callback failure {}: {}",
job.id, mw::errorMsg(failure.error()));
}
active_job.reset();
}
catch(const std::exception& error)
{
spdlog::error("Callback worker failed: {}", error.what());
if(active_job.has_value())
{
auto failure = store_.fail(active_job->id, "Callback worker failed",
true, std::nullopt);
if(!failure.has_value())
{
spdlog::error("Unable to recover callback job {}: {}",
active_job->id, mw::errorMsg(failure.error()));
}
}
std::unique_lock lock(mutex_);
condition_.wait_for(lock, stop_token, std::chrono::seconds(1),
[] { return false; });
}
catch(...)
{
spdlog::error("Callback worker failed with an unknown exception");
if(active_job.has_value())
{
auto failure = store_.fail(active_job->id, "Callback worker failed",
true, std::nullopt);
if(!failure.has_value())
{
spdlog::error("Unable to recover callback job {}: {}",
active_job->id, mw::errorMsg(failure.error()));
}
}
std::unique_lock lock(mutex_);
condition_.wait_for(lock, stop_token, std::chrono::seconds(1),
[] { return false; });
}
}
}
} // namespace telegrammer