2020-05-20 09:40:49 +00:00
|
|
|
#pragma once
|
|
|
|
|
2020-06-02 13:15:53 +00:00
|
|
|
#include <thread>
|
2020-05-20 09:40:49 +00:00
|
|
|
#include <memory>
|
2020-06-02 13:15:53 +00:00
|
|
|
#include <mutex>
|
2020-05-20 09:40:49 +00:00
|
|
|
#include <amqpcpp.h>
|
|
|
|
#include <amqpcpp/libevent.h>
|
|
|
|
#include <amqpcpp/linux_tcp.h>
|
|
|
|
#include <common/types.h>
|
|
|
|
#include <event2/event.h>
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
class RabbitMQHandler : public AMQP::LibEventHandler
|
|
|
|
{
|
|
|
|
|
|
|
|
public:
|
|
|
|
RabbitMQHandler(event_base * evbase_, Poco::Logger * log_);
|
|
|
|
|
|
|
|
void onError(AMQP::TcpConnection * connection, const char * message) override;
|
2020-06-09 21:52:06 +00:00
|
|
|
void startConsumerLoop(std::atomic<bool> & loop_started);
|
2020-06-07 11:14:05 +00:00
|
|
|
void startProducerLoop();
|
|
|
|
void stopWithTimeout();
|
2020-05-20 09:40:49 +00:00
|
|
|
void stop();
|
2020-06-11 20:05:35 +00:00
|
|
|
std::atomic<bool> & checkStopIsScheduled() { return stop_scheduled; };
|
2020-05-20 09:40:49 +00:00
|
|
|
|
|
|
|
private:
|
|
|
|
event_base * evbase;
|
|
|
|
Poco::Logger * log;
|
2020-05-29 16:04:44 +00:00
|
|
|
|
2020-06-07 11:14:05 +00:00
|
|
|
timeval tv;
|
2020-06-11 20:05:35 +00:00
|
|
|
std::atomic<bool> stop_scheduled = false;
|
2020-06-04 06:22:53 +00:00
|
|
|
std::timed_mutex mutex_before_event_loop;
|
2020-06-11 20:05:35 +00:00
|
|
|
std::mutex mutex_before_loop_stop;
|
2020-05-20 09:40:49 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
}
|