#pragma once
#include <condition_variable>
#include <cstddef>
#include <chrono>
#include <mutex>
#include <thread>
#include <vector>
#include "delivery_store.hpp"
namespace telegrammer
{
/// Delivers queued callbacks with bounded concurrency and retries.
class Dispatcher
{
public:
/// Construct a dispatcher with a fixed number of callback workers.
Dispatcher(DeliveryStore& store, std::size_t worker_count = 4);
/// Start callback workers.
void start();
/// Request workers to stop and wait for active callbacks.
void stop();
/// Wake workers after new delivery jobs are committed.
void notify();
private:
DeliveryStore& store_;
std::size_t worker_count_;
std::mutex mutex_;
std::condition_variable_any condition_;
std::mutex purge_mutex_;
std::chrono::steady_clock::time_point next_purge_{};
std::vector<std::jthread> workers_;
bool shouldPurge();
void run(std::stop_token stop_token);
};
} // namespace telegrammer