Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 34 additions & 21 deletions include/DataSender.h
Original file line number Diff line number Diff line change
Expand Up @@ -12,26 +12,34 @@
#include <Buffer.h>
#include <DAQHeader.h>
#include <functional>
#include <fstream>
#include <filesystem>
#include <string>

using namespace ToolFramework;

void DeleteHeader(void* data);

// Disk overflow cache (see DataSender.cpp): caches to cache_file instead of
// dropping data past retry_limit. max_size_bytes caps its size (0 = unbounded).
bool CacheMessageToDisk(const std::string& cache_file, DataMessages* msg, uint64_t max_size_bytes = 0);
DataMessages* RecoverMessageFromDisk(const std::string& cache_file);

struct DataSender_args: Thread_args{

DataSender_args();
~DataSender_args();

DAQInterface* daq_interface;
Buffer<DataMessages*>* in_buffer;
std::vector<DataMessages*> new_send;
std::deque<DataMessages*> to_send;
std::vector<DataMessages*> new_send;
std::deque<DataMessages*> to_send;
std::map<uint32_t, DataMessages*> sent;
std::chrono::milliseconds time_span;
std::map<uint32_t, DataMessages*>::iterator message_it;

uint16_t retry_limit = 200;
uint32_t resend_period_ms = 50;
uint32_t resend_period_ms = 50;
uint32_t poll_period_ms = 100;
uint16_t send_error_counter = 0;

Expand All @@ -42,31 +50,36 @@ struct DataSender_args: Thread_args{
std::atomic<uint64_t> num_data_messages;
std::atomic<uint64_t> num_data_akn;
std::atomic<uint64_t> num_data_deleted;


std::string cache_file = ""; // empty = disk overflow cache disabled (old behaviour: drop + log)
uint32_t cache_recover_per_cycle = 1;
uint64_t cache_max_size_bytes = 0; // 0 = unbounded; cap on cache_file's size
std::atomic<uint64_t> num_data_cached{0}; // current number of messages sitting in cache_file

};

class DataSender{

public:

DataSender(DAQInterface* interface, std::string config_file);
~DataSender();

bool LoadConfig(std::string json);
bool LoadConfig(Store& vars);
bool Add(void* data, size_t size, uint32_t coarse_counter, std::function<void(void*)> del_func);
bool Add(DataMessages* message);

std::string Summary();



private:

static void Thread(Thread_args* arg);

DAQInterface* daq_interface;
Utilities m_utils;
Utilities m_utils;
DataSender_args args;
Buffer<DataMessages*> in_buffer;

Expand All @@ -82,22 +95,22 @@ class DataSender{
int32_t tcp_keepalive_idle_sec = 5;
int32_t tcp_keepalive_count = 12;
int32_t tcp_keepalive_interval_sec = 5;

std::string data_port = "";

uint64_t message_num = 0;

uint64_t num_data_messages = 0;
uint64_t num_data_akn = 0;
uint64_t num_data_deleted = 0;
float num_data_messages_rate = 0;
float num_data_akn_rate = 0;
float num_data_messages_rate = 0;
float num_data_akn_rate = 0;
float num_data_deleted_rate = 0;
std::chrono::time_point<std::chrono::steady_clock> last;

uint8_t card_type = 0;
uint16_t card_id = 0;

};

#endif
Loading