BareGit
#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