From 4746e0afa18b7172d585efe0906075acd490b4d8 Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Mon, 10 Aug 2026 12:19:47 +0100 Subject: [PATCH 1/8] add compression and decompression of slow controls and alert payloads --- .../SlowControlCollection.cpp | 349 +++++++++++------- src/ServiceDiscovery/SlowControlCollection.h | 12 + 2 files changed, 221 insertions(+), 140 deletions(-) diff --git a/src/ServiceDiscovery/SlowControlCollection.cpp b/src/ServiceDiscovery/SlowControlCollection.cpp index 13e24e4..1765f2e 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -2,6 +2,10 @@ using namespace ToolFramework; +namespace { + const unsigned char ZSTD_MAGIC_BYTES[4] = {0x28,0xB5,0x2F,0xFD}; // ZSTD_MAGICNUMBER from zstd.h BUT REVERSED! +} + SlowControlCollectionThread_args::SlowControlCollectionThread_args(){ sock=0; @@ -12,9 +16,6 @@ SlowControlCollectionThread_args::SlowControlCollectionThread_args(){ alert_functions_mutex=0; SC_vars=0; - m_pub = 0; - pub_monitor_socket = 0; - pub_connected_mtx = 0; } @@ -55,6 +56,9 @@ SlowControlCollection::SlowControlCollection(){ Add("NewConfig",SlowControlElementType(INFO),0,0,false,false); SC_vars["NewConfig"]->SetValue(0); + zstd_cctx = ZSTD_createCCtx(); + zstd_dctx = ZSTD_createDCtx(); + } SlowControlCollection::~SlowControlCollection(){ @@ -123,15 +127,15 @@ bool SlowControlCollection::Init(zmq::context_t* context, int sc_port, bool new_ if(!m_util->AddPort("alerts",alert_send_port)){ - delete m_pub; - m_pub=0; - m_alerts_send = false; + delete m_pub; + m_pub=0; + m_alerts_send = false; - delete args; - args=0; + delete args; + args=0; - std::clog<<"Error adding alert send port to SD"<AddPort("alertr",alert_receive_port)){ - delete args->sub; - args->sub = 0; - m_alerts_receive = false; + delete args->sub; + args->sub = 0; + m_alerts_receive = false; - delete args; - args=0; + delete args; + args=0; - - std::clog<<"Error adding port alert receive to SD"<(message.data())); + std::string payload; + if(!args->SCC->ZstdDecompress(args->SCC, (char*)message.data(), message.size(), payload)){ + std::cerr<<"failed to decompress slow control message!"<SCC->Print()=%s\n", args->SCC->Print().c_str()); if(value=="JSON"){ - reply=args->SCC->PrintJSON(); - strip=true; + reply=args->SCC->PrintJSON(); + strip=true; } else reply=args->SCC->Print(); //printf("reply=%s\n", reply.c_str()); @@ -336,29 +344,29 @@ void SlowControlCollection::Thread(Thread_args* arg){ } else{ - reply=key; - if((*args->SCC)[key]->GetType() == SlowControlElementType(BUTTON)){ - value="1"; - } + reply=key; + if((*args->SCC)[key]->GetType() == SlowControlElementType(BUTTON)){ + value="1"; + } //std::stringstream input; //input<("msg_value"); - // std::string key=""; - // std::string value=""; + // std::string key=""; + // std::string value=""; //input>>key>>value; - //printf("d0 %s = %s : %s\n", reply.c_str(), key.c_str(), value.c_str()); - if(value!=""){ - (*args->SCC)[key]->SetValue(value); - //(*args->SCC)[key]->Print(); - SCFunction tmp_func= (*args->SCC)[key]->GetChangeFunction(); - if (tmp_func!=nullptr) reply=tmp_func(key.c_str()); - - } - else{ - SCFunction tmp_func= (*args->SCC)[key]->GetReadFunction(); - if (tmp_func!=nullptr) reply=tmp_func(key.c_str()); - else (*args->SCC)[key]->GetValue(reply); - - } + //printf("d0 %s = %s : %s\n", reply.c_str(), key.c_str(), value.c_str()); + if(value!=""){ + (*args->SCC)[key]->SetValue(value); + //(*args->SCC)[key]->Print(); + SCFunction tmp_func= (*args->SCC)[key]->GetChangeFunction(); + if (tmp_func!=nullptr) reply=tmp_func(key.c_str()); + + } + else{ + SCFunction tmp_func= (*args->SCC)[key]->GetReadFunction(); + if (tmp_func!=nullptr) reply=tmp_func(key.c_str()); + else (*args->SCC)[key]->GetValue(reply); + + } } } */ @@ -372,9 +380,8 @@ void SlowControlCollection::Thread(Thread_args* arg){ std::string tmp2=""; rr>>tmp2; //printf("reply message is= %s \n",tmp2.c_str()); - zmq::message_t send(tmp2.length()+1); - snprintf ((char *) send.data(), tmp2.length()+1 , "%s" ,tmp2.c_str()) ; + zmq::message_t send = args->SCC->ZstdCompress(args->SCC, tmp2); bool tmp_ok = args->sock->send(identity, ZMQ_SNDMORE); if(tmp_ok) tmp_ok = tmp_ok && args->sock->send(blank, ZMQ_SNDMORE); @@ -407,17 +414,20 @@ void SlowControlCollection::Thread(Thread_args* arg){ ok = args->sub->recv(&message); if(ok==0){ // FIXME this case should be handled! what do we do? - std::cerr<<"failed to receive alert payload!"<SCC->ZstdDecompress(args->SCC, (char*)message.data(), message.size(), payload)){ + std::cerr<<"failed to decompress "<sub->recv(&message); + args->sub->recv(&message); // FIXME do we want any warnings or handling here? //memcpy((void*)payload.data(),message.data(),message.size()); //a++; } @@ -427,8 +437,8 @@ void SlowControlCollection::Thread(Thread_args* arg){ if(iss.str() == "LoadConfig") (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadStart); else if(iss.str() == "ChangeConfig"){ if((*args->SC_vars)["NewConfig"]->GetValue() == 0){ - args->alert_functions_mutex->unlock(); - return; + args->alert_functions_mutex->unlock(); + return; } (*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeStart); } @@ -438,30 +448,30 @@ void SlowControlCollection::Thread(Thread_args* arg){ if(args->alert_functions->count(iss.str())){ if(has_data){ - try{ - error = !((*(args->alert_functions))[iss.str()](iss.str().c_str(), payload.c_str())); - } - catch(...){ - error = true; - } - if(iss.str() == "LoadConfig"){ - if(error)(*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadFail); - else (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadEnd); - } - + try{ + error = !((*(args->alert_functions))[iss.str()](iss.str().c_str(), payload.c_str())); + } + catch(...){ + error = true; + } + if(iss.str() == "LoadConfig"){ + if(error)(*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadFail); + else (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadEnd); + } + } else { - try{ - error=!((*(args->alert_functions))[iss.str()](iss.str().c_str(), 0)); - } - catch(...){ - error = true; - } - if(iss.str() == "ChangeConfig"){ - if(error)(*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeFail); - else (*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeEnd); - (*args->SC_vars)["NewConfig"]->SetValue(0); - } + try{ + error=!((*(args->alert_functions))[iss.str()](iss.str().c_str(), 0)); + } + catch(...){ + error = true; + } + if(iss.str() == "ChangeConfig"){ + if(error)(*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeFail); + else (*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeEnd); + (*args->SC_vars)["NewConfig"]->SetValue(0); + } } if(error) std::cerr<<"alert fucntion failed: "<send(message, ZMQ_SNDMORE); if(!ok) return false; // err: "zmq send "+zmq_strerror(errno) - zmq::message_t message2(payload.length()+1); - snprintf((char*) message2.data(), payload.length()+1, "%s", payload.c_str()); + + zmq::message_t message2 = args->SCC->ZstdCompress(args->SCC, payload); return m_pub->send(message2); // err: "zmq send "+zmq_strerror(errno) } @@ -644,16 +654,16 @@ void SlowControlCollection::Unpack(std::string in, std::mapsecond.length(); i++){ - if(it->second[i]=='{') first=i; - if(it->second[i]=='}'){ - // std::stringstream tmp; - //tmp<second.substr(first,i-first+1)); - out[header+tmp.Get("name")]=it->second.substr(first,i-first+1); - // counter++; - } - + if(it->second[i]=='{') first=i; + if(it->second[i]=='}'){ + // std::stringstream tmp; + //tmp<second.substr(first,i-first+1)); + out[header+tmp.Get("name")]=it->second.substr(first,i-first+1); + // counter++; + } + } } @@ -665,17 +675,17 @@ void SlowControlCollection::Unpack(std::string in, std::mapsecond.length(); i++){ - //std::cout<<"i="<second[i]="<second[i]<<" : counter="<second[i]=='}') bracket_counter--; + //std::cout<<"i="<second[i]="<second[i]<<" : counter="<second[i]=='}') bracket_counter--; } } } @@ -709,8 +719,8 @@ bool SlowControlCollection::Update(SlowControlCollection* SCC, std::string key, //std::cout<<"variable exists"<GetType() == SlowControlElementType(INFO)){ if(!(*SCC)[key]->GetValue(value)){ - reply="Error getting value form key: "+key; - return false; + reply="Error getting value form key: "+key; + return false; } reply=value; return true; @@ -719,7 +729,7 @@ bool SlowControlCollection::Update(SlowControlCollection* SCC, std::string key, else{ reply=key; if((*SCC)[key]->GetType() == SlowControlElementType(BUTTON)){ - value="1"; + value="1"; } //std::stringstream input; //input<("msg_value"); @@ -728,49 +738,49 @@ bool SlowControlCollection::Update(SlowControlCollection* SCC, std::string key, //input>>key>>value; //printf("d0 %s = %s : %s\n", reply.c_str(), key.c_str(), value.c_str()); if(value!=""){ - if(!testing || (testing && !(*SCC)[key]->Lockable())){ - - if(!(*SCC)[key]->SetValue(value)){ - reply =" Error setting "+key+" to value: " + value; - return false; - } - else{ - reply = value; - return true; - } - //(*SCC)[key]->Print(); - /* - SCFunction tmp_func= (*SCC)[key]->GetChangeFunction(); - if (tmp_func!=nullptr){ - try{ - reply=tmp_func(key.c_str()); - - } - catch(...){ - reply= "change function failed"; - } - } - */ - } - else reply = key + " locked"; + if(!testing || (testing && !(*SCC)[key]->Lockable())){ + + if(!(*SCC)[key]->SetValue(value)){ + reply =" Error setting "+key+" to value: " + value; + return false; + } + else{ + reply = value; + return true; + } + //(*SCC)[key]->Print(); + /* + SCFunction tmp_func= (*SCC)[key]->GetChangeFunction(); + if (tmp_func!=nullptr){ + try{ + reply=tmp_func(key.c_str()); + + } + catch(...){ + reply= "change function failed"; + } + } + */ + } + else reply = key + " locked"; } else{ - /* - SCFunction tmp_func= (*SCC)[key]->GetReadFunction(); - if (tmp_func!=nullptr){ - try{ - reply=tmp_func(key.c_str()); - } - catch(...){ - reply="read function failed"; - } - } - else (*SCC)[key]->GetValue(reply); - */ + /* + SCFunction tmp_func= (*SCC)[key]->GetReadFunction(); + if (tmp_func!=nullptr){ + try{ + reply=tmp_func(key.c_str()); + } + catch(...){ + reply="read function failed"; + } + } + else (*SCC)[key]->GetValue(reply); + */ if(!(*SCC)[key]->GetValue(reply)){ - reply="Error getting value from key: "+key; - return false; - } + reply="Error getting value from key: "+key; + return false; + } } } return true; @@ -872,3 +882,62 @@ void SlowControlCollection::ClearState(){ SC_vars["State"]->SetValue(m_state); return; } + +zmq::message_t SlowControlCollection::ZstdCompress(SlowControlCollection* SCC, std::string& msg){ + if(msg.length()COMPRESS_THRESHOLD){ + zmq::message_t zmsg(msg.size()); + memcpy(zmsg.data(), msg.data(), msg.size()); + return zmsg; + } + + std::unique_lock locker(*SCC->zstd_cctx_mtx); + std::string compressed_msg_buf; + compressed_msg_buf.resize(ZSTD_compressBound(msg.size())); + uint64_t bytes_to_send = ZSTD_compressCCtx(SCC->zstd_cctx, (void*)compressed_msg_buf.data(), compressed_msg_buf.size(), msg.data(), msg.size(), SCC->zstd_compression_level); + if(ZSTD_isError(bytes_to_send)){ + locker.unlock(); + std::string errmsg = std::string{"Warning: error compressing multicast message "}+ZSTD_getErrorName(bytes_to_send); + std::clog << errmsg << std::endl; + // send it uncompressed + zmq::message_t zmsg(msg.size()); + memcpy(zmsg.data(), msg.data(), msg.size()); + return zmsg; + } + zmq::message_t zmsg(msg.size()); + memcpy(zmsg.data(), compressed_msg_buf.data(), bytes_to_send); + return zmsg; +} + +bool SlowControlCollection::ZstdDecompress(SlowControlCollection* SCC, char* msg, uint64_t msgsize, std::string& decompress_buffer){ + std::string errmsg; + std::string* decompressed_msg=nullptr; + std::unique_lock locker(*SCC->zstd_dctx_mtx); + if(msgsize>4 && std::memcmp(msg,ZSTD_MAGIC_BYTES,4)==0){ + uint64_t decompressed_bytes = ZSTD_getFrameContentSize(msg, msgsize); + if(decompressed_bytes==ZSTD_CONTENTSIZE_UNKNOWN || decompressed_bytes==ZSTD_CONTENTSIZE_ERROR){ + // bad response + errmsg = std::string{"Received corrupt zstd message "}+ZSTD_getErrorName(decompressed_bytes); + goto decompress_error; + } + if(decompressed_bytes > SCC->MAX_DECOMPRESSED_SIZE){ + errmsg = "Compressed message with oversized payload: "+std::to_string(decompressed_bytes)+" bytes"; + goto decompress_error; + } + decompress_buffer.resize(decompressed_bytes); + decompressed_bytes = ZSTD_decompressDCtx(SCC->zstd_dctx,(void*)decompress_buffer.data(),decompressed_bytes, msg, msgsize); + if(ZSTD_isError(decompressed_bytes)){ + errmsg = std::string{"zstd error decompressing response: "}+ZSTD_getErrorName(decompressed_bytes); + goto decompress_error; + } + } else { + // message not compressed + decompress_buffer.assign(msg, msgsize); + } + return true; + + decompress_error: + locker.unlock(); + std::clog << errmsg << std::endl; + decompress_buffer.clear(); + return false; +} diff --git a/src/ServiceDiscovery/SlowControlCollection.h b/src/ServiceDiscovery/SlowControlCollection.h index c2a43b1..fb3abce 100644 --- a/src/ServiceDiscovery/SlowControlCollection.h +++ b/src/ServiceDiscovery/SlowControlCollection.h @@ -5,6 +5,8 @@ #include #include "DAQUtilities.h" #include +#include +#include namespace ToolFramework{ @@ -68,6 +70,8 @@ namespace ToolFramework{ void SetError(bool error); void SetWarning(bool warn); void ClearState(); + zmq::message_t ZstdCompress(SlowControlCollection* SCC, std::string& msg); + bool ZstdDecompress(SlowControlCollection* SCC, char* msg, uint64_t msgsize, std::string& decompress_buffer); template T GetValue(std::string name){ if(!SC_vars.count(name)) return T{}; @@ -81,6 +85,14 @@ namespace ToolFramework{ std::map m_alert_functions; std::mutex m_alert_functions_mutex; + ZSTD_CCtx* zstd_cctx; + std::mutex* zstd_cctx_mtx; + ZSTD_DCtx* zstd_dctx; + std::mutex* zstd_dctx_mtx; + int zstd_compression_level=1; + uint32_t COMPRESS_THRESHOLD=0; //1024; // compress any send messages > this many bytes + uint32_t MAX_DECOMPRESSED_SIZE=655355; // refuse to decompress messages that will exceed this size once decompressed + DAQUtilities* m_util; zmq::context_t* m_context; zmq::socket_t* m_pub; From a03d55472f01a5deda9079bd8ad585de0b4567f6 Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Mon, 10 Aug 2026 12:45:09 +0100 Subject: [PATCH 2/8] add decompression to RemoteControl --- src/RemoteControl/RemoteControl.cpp | 69 +++++++++++++++---- .../SlowControlCollection.cpp | 1 - 2 files changed, 56 insertions(+), 14 deletions(-) diff --git a/src/RemoteControl/RemoteControl.cpp b/src/RemoteControl/RemoteControl.cpp index 20a7ddc..a8a6ab7 100644 --- a/src/RemoteControl/RemoteControl.cpp +++ b/src/RemoteControl/RemoteControl.cpp @@ -5,6 +5,7 @@ #include "ServiceDiscovery.h" #include "zmq.hpp" +#include #include // uuid class #include // generators @@ -18,9 +19,13 @@ #define GROUP_COMMAND_REPLY_WAIT 2000 #define FILE_SEND_WAIT 120000 #define FILE_SEND_PORT 24001 +#define MAX_DECOMPRESSED_SIZE 655355 +const unsigned char ZSTD_MAGIC_BYTES[4] = {0x28,0xB5,0x2F,0xFD}; // ZSTD_MAGICNUMBER from zstd.h BUT REVERSED! using namespace ToolFramework; +bool ZstdDecompress(ZSTD_DCtx* zstd_dctx, char* msg, uint64_t msgsize, std::string& decompress_buffer); + int main(int argc, char** argv){ // if (argc!=3) return 1; @@ -30,6 +35,7 @@ int main(int argc, char** argv){ zmq::context_t context(3); + ZSTD_DCtx* zstd_dctx = ZSTD_createDCtx(); //std::string address(argv[1]); // std::stringstream tmp (argv[2]); @@ -236,14 +242,17 @@ int main(int argc, char** argv){ zmq::message_t receive; if(ServiceSend.recv(&receive)){ - std::istringstream iss(static_cast(receive.data())); - std::string answer; - answer=iss.str(); + if(!ZstdDecompress(zstd_dctx, (char*)receive.data(), receive.size(), answer)){ + std::cerr<<"failed to decompress reply!"<("msg_type")=="Command Reply") std::cout<("msg_value")<("msg_type")=="Command Reply") std::cout<("msg_value"))<(receive.data())); - std::string answer; - answer=iss.str(); - - Store rr; - rr.JsonParser(answer); - if(rr.Get("msg_type")=="Command Reply") std::cout<("msg_value")<("msg_type")=="Command Reply") std::cout<("msg_value")<4 && std::memcmp(msg,ZSTD_MAGIC_BYTES,4)==0){ + uint64_t decompressed_bytes = ZSTD_getFrameContentSize(msg, msgsize); + if(decompressed_bytes==ZSTD_CONTENTSIZE_UNKNOWN || decompressed_bytes==ZSTD_CONTENTSIZE_ERROR){ + // bad response + errmsg = std::string{"Received corrupt zstd message "}+ZSTD_getErrorName(decompressed_bytes); + goto decompress_error; + } + if(decompressed_bytes > MAX_DECOMPRESSED_SIZE){ + errmsg = "Compressed message with oversized payload: "+std::to_string(decompressed_bytes)+" bytes"; + goto decompress_error; + } + decompress_buffer.resize(decompressed_bytes); + decompressed_bytes = ZSTD_decompressDCtx(zstd_dctx,(void*)decompress_buffer.data(),decompressed_bytes, msg, msgsize); + if(ZSTD_isError(decompressed_bytes)){ + errmsg = std::string{"zstd error decompressing response: "}+ZSTD_getErrorName(decompressed_bytes); + goto decompress_error; + } + } else { + // message not compressed + decompress_buffer.assign(msg, msgsize); + } + return true; + + decompress_error: + std::cerr << errmsg << std::endl; + decompress_buffer.clear(); + return false; +} diff --git a/src/ServiceDiscovery/SlowControlCollection.cpp b/src/ServiceDiscovery/SlowControlCollection.cpp index 1765f2e..8bf84f2 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -910,7 +910,6 @@ zmq::message_t SlowControlCollection::ZstdCompress(SlowControlCollection* SCC, s bool SlowControlCollection::ZstdDecompress(SlowControlCollection* SCC, char* msg, uint64_t msgsize, std::string& decompress_buffer){ std::string errmsg; - std::string* decompressed_msg=nullptr; std::unique_lock locker(*SCC->zstd_dctx_mtx); if(msgsize>4 && std::memcmp(msg,ZSTD_MAGIC_BYTES,4)==0){ uint64_t decompressed_bytes = ZSTD_getFrameContentSize(msg, msgsize); From f02092a102e6da60501e191738f3e115792729ec Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Wed, 19 Aug 2026 11:39:57 +0100 Subject: [PATCH 3/8] slow control compression bugfixes --- src/ServiceDiscovery/SlowControlCollection.cpp | 4 ++-- src/ServiceDiscovery/SlowControlCollection.h | 6 +++--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/src/ServiceDiscovery/SlowControlCollection.cpp b/src/ServiceDiscovery/SlowControlCollection.cpp index 8bf84f2..11e8a7d 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -890,7 +890,7 @@ zmq::message_t SlowControlCollection::ZstdCompress(SlowControlCollection* SCC, s return zmsg; } - std::unique_lock locker(*SCC->zstd_cctx_mtx); + std::unique_lock locker(SCC->zstd_cctx_mtx); std::string compressed_msg_buf; compressed_msg_buf.resize(ZSTD_compressBound(msg.size())); uint64_t bytes_to_send = ZSTD_compressCCtx(SCC->zstd_cctx, (void*)compressed_msg_buf.data(), compressed_msg_buf.size(), msg.data(), msg.size(), SCC->zstd_compression_level); @@ -910,7 +910,7 @@ zmq::message_t SlowControlCollection::ZstdCompress(SlowControlCollection* SCC, s bool SlowControlCollection::ZstdDecompress(SlowControlCollection* SCC, char* msg, uint64_t msgsize, std::string& decompress_buffer){ std::string errmsg; - std::unique_lock locker(*SCC->zstd_dctx_mtx); + std::unique_lock locker(SCC->zstd_dctx_mtx); if(msgsize>4 && std::memcmp(msg,ZSTD_MAGIC_BYTES,4)==0){ uint64_t decompressed_bytes = ZSTD_getFrameContentSize(msg, msgsize); if(decompressed_bytes==ZSTD_CONTENTSIZE_UNKNOWN || decompressed_bytes==ZSTD_CONTENTSIZE_ERROR){ diff --git a/src/ServiceDiscovery/SlowControlCollection.h b/src/ServiceDiscovery/SlowControlCollection.h index fb3abce..82ddd08 100644 --- a/src/ServiceDiscovery/SlowControlCollection.h +++ b/src/ServiceDiscovery/SlowControlCollection.h @@ -86,11 +86,11 @@ namespace ToolFramework{ std::mutex m_alert_functions_mutex; ZSTD_CCtx* zstd_cctx; - std::mutex* zstd_cctx_mtx; + std::mutex zstd_cctx_mtx; ZSTD_DCtx* zstd_dctx; - std::mutex* zstd_dctx_mtx; + std::mutex zstd_dctx_mtx; int zstd_compression_level=1; - uint32_t COMPRESS_THRESHOLD=0; //1024; // compress any send messages > this many bytes + uint32_t COMPRESS_THRESHOLD=1024000000; // compress any send messages > this many bytes uint32_t MAX_DECOMPRESSED_SIZE=655355; // refuse to decompress messages that will exceed this size once decompressed DAQUtilities* m_util; From 9f1eeb5bacd68fa0c451b2f9c9a0f3dc6b907f5c Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Thu, 20 Aug 2026 00:38:45 +0100 Subject: [PATCH 4/8] move multicast max size check to after compression, and reduce to account for IP header --- src/ServiceDiscovery/ServicesBackend.cpp | 43 ++++++++++---------- src/ServiceDiscovery/SlowControlCollection.h | 5 ++- 2 files changed, 25 insertions(+), 23 deletions(-) diff --git a/src/ServiceDiscovery/ServicesBackend.cpp b/src/ServiceDiscovery/ServicesBackend.cpp index 2c50868..0159f1a 100644 --- a/src/ServiceDiscovery/ServicesBackend.cpp +++ b/src/ServiceDiscovery/ServicesBackend.cpp @@ -1,8 +1,8 @@ #include "ServicesBackend.h" namespace { - const uint32_t MAX_UDP_PACKET_SIZE = 655355; - const uint32_t MAX_DECOMPRESSED_MSG_SIZE = 655355; + const uint32_t MAX_UDP_PACKET_SIZE = 65507; // limit from UDP message length field + const uint32_t MAX_DECOMPRESSED_MSG_SIZE = 104857600; // 100MB limit. I'm sure we can spare that much RAM. const unsigned char ZSTD_MAGIC_BYTES[4] = {0x28,0xB5,0x2F,0xFD}; // ZSTD_MAGICNUMBER from zstd.h BUT REVERSED! } @@ -142,7 +142,7 @@ bool ServicesBackend::Initialise(Store &variables_in){ if(msg_compression){ zstd_cctx = ZSTD_createCCtx(); - compressed_msg_buf = new char[ZSTD_compressBound(MAX_UDP_PACKET_SIZE)]; + compressed_msg_buf = new char[MAX_UDP_PACKET_SIZE]; zstd_dctx = ZSTD_createDCtx(); } @@ -433,24 +433,25 @@ bool ServicesBackend::SendMulticast(MulticastType type, std::string command, std bytes_to_send = command.length(); } - /* - // check for listeners...? - seems redundant, multicast can always send - zmq::poll(&multicast_poller,1, 0); // timeout 0 = return immediately... - if(multicast_poller.revents & ZMQ_POLLOUT){ - */ - - // got a listener - ship it - socket_mtx->lock(); - int cnt = sendto(multicast_socket, msg_to_send, bytes_to_send, 0, (struct sockaddr*)multicast_addr, multicast_addrlen); - socket_mtx->unlock(); - if(cnt < 0){ - std::string errmsg = "Error sending multicast message: "+std::string{strerror(errno)}; - Log(errmsg,v_error,m_verbosity); - if(err) *err= errmsg; //zmq_strerror(errno); - return false; - } - - //} + // check we're not going to exceed multicast message size limit + if(bytes_to_send > MAX_UDP_PACKET_SIZE){ + // we can't send this on multicast. + if(locker.owns_lock()) locker.unlock(); + std::string errmsg = "Error: message exceeds maximum bytes of "+std::to_string(MAX_UDP_PACKET_SIZE); + Log(errmsg,v_error,m_verbosity); // XXX should send to MM uncompressed, along with other errors + if(err) *err= errmsg; + return false; + } + + socket_mtx->lock(); + int cnt = sendto(multicast_socket, msg_to_send, bytes_to_send, 0, (struct sockaddr*)multicast_addr, multicast_addrlen); + socket_mtx->unlock(); + if(cnt < 0){ + std::string errmsg = "Error sending multicast message: "+std::string{strerror(errno)}; + Log(errmsg,v_error,m_verbosity); + if(err) *err= errmsg; //zmq_strerror(errno); + return false; + } return true; } diff --git a/src/ServiceDiscovery/SlowControlCollection.h b/src/ServiceDiscovery/SlowControlCollection.h index 82ddd08..f189d3b 100644 --- a/src/ServiceDiscovery/SlowControlCollection.h +++ b/src/ServiceDiscovery/SlowControlCollection.h @@ -90,8 +90,9 @@ namespace ToolFramework{ ZSTD_DCtx* zstd_dctx; std::mutex zstd_dctx_mtx; int zstd_compression_level=1; - uint32_t COMPRESS_THRESHOLD=1024000000; // compress any send messages > this many bytes - uint32_t MAX_DECOMPRESSED_SIZE=655355; // refuse to decompress messages that will exceed this size once decompressed + uint32_t COMPRESS_THRESHOLD=1024; // compress any send messages > this many bytes + uint32_t MAX_DECOMPRESSED_SIZE=104857600; // refuse to decompress messages that will exceed this size once decompressed + // defautl 100MB; I'm sure we can spare that much RAM. DAQUtilities* m_util; zmq::context_t* m_context; From 8e9f4895d4e8f4cedafdd7508033a3d5ace06cf7 Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Mon, 24 Aug 2026 10:53:11 +0100 Subject: [PATCH 5/8] * make LocalConfig slow control hidden until web page supports showing JSON * send logging and monitoring batches on receipt of a new message when the new message would push the batch size over the MTU. * reject logging and montioring messages that would exceed maximum UDP payload size after compression * increase verbosity level of some prints in ServicesBackend * Make SlowControlCollection use std::unique_lock instead of manual locking in Thread * bugfix: AlertSubscribe should check alert receiving is enabled, not sending * bugfix: SlowControlCollection response compression size was incorrect. * add error checking when determining SlowControlCollection compressed response size * reduce compression threshold to 50 bytes --- configfiles/template/ToolChainConfig | 1 + src/ServiceDiscovery/Services.cpp | 208 +++++++++++++----- src/ServiceDiscovery/Services.h | 7 + src/ServiceDiscovery/ServicesBackend.cpp | 13 +- .../SlowControlCollection.cpp | 26 +-- src/ServiceDiscovery/SlowControlCollection.h | 2 +- 6 files changed, 175 insertions(+), 82 deletions(-) diff --git a/configfiles/template/ToolChainConfig b/configfiles/template/ToolChainConfig index b4dd22b..21a9768 100644 --- a/configfiles/template/ToolChainConfig +++ b/configfiles/template/ToolChainConfig @@ -45,6 +45,7 @@ mon_port 5000 # remote multicast port to send monitoring messa multicast_send_period 5000 # logging & monitoring messages will be sent as a batch once per this period alarm_cooldown_ms 1000 # alarms with the same device and message will be limited to this rate mon_merge_period 5000 # monitoring messages with the same device and subject will be limited to this rate +MTU 1500 # MTU of interface used for logging/monitoring ##### Tools To Add ##### Tools_File configfiles/ToolsConfig # list of tools to run and their config files diff --git a/src/ServiceDiscovery/Services.cpp b/src/ServiceDiscovery/Services.cpp index cec329d..c19f7b6 100644 --- a/src/ServiceDiscovery/Services.cpp +++ b/src/ServiceDiscovery/Services.cpp @@ -3,8 +3,9 @@ using namespace ToolFramework; namespace { - constexpr uint32_t MAX_UDP_PACKET_SIZE = 655355; + constexpr uint32_t MAX_UDP_PACKET_SIZE = 65507; constexpr size_t MAX_MSG_SIZE = MAX_UDP_PACKET_SIZE-100; // 100 chars for JSON keys, timestamp string and quotes/commas + uint32_t MSS_SIZE=1472; // for standard MTU of 1500 bytes } Services::Services(){ @@ -55,6 +56,7 @@ bool Services::Init(Store &m_variables, zmq::context_t* context_in, SlowControlC m_variables.Get("multicast_send_period_ms",multicast_send_period_ms); m_variables.Get("alarm_cooldown_ms",alarm_cooldown_ms); m_variables.Get("verbose",m_verbose); + if(m_variables.Get("MTU",MSS_SIZE)) MSS_SIZE -= 28; // account for headers sc_vars->InitThreadedReceiver(m_context, sc_port, 100, new_service, alert_receive_port, alerts_receive, alert_send_port, alerts_send); m_backend_client.SetUp(m_context); @@ -65,7 +67,7 @@ bool Services::Init(Store &m_variables, zmq::context_t* context_in, SlowControlC sc_vars->Add("LoadConfig",SlowControlElementType(COMMAND),std::bind(&Services::LoadConfigSlowControlFunc, this, std::placeholders::_1),0,false,false); AlertSubscribe("LoadConfig", std::bind(&Services::LoadConfigAlertFunc, this, std::placeholders::_1, std::placeholders::_2)); - sc_vars->Add("LocalConfig",SlowControlElementType(INFO),std::bind(&Services::SCLocalConfig, this, std::placeholders::_1),0,false,false); + sc_vars->Add("LocalConfig",SlowControlElementType(INFO),0,std::bind(&Services::SCLocalConfig, this, std::placeholders::_1),false,true); // FIXME hidden until Control page supports JSON if(!m_variables.Get("service_name",m_name)) m_name="test_service"; @@ -106,6 +108,13 @@ bool Services::Init(Store &m_variables, zmq::context_t* context_in, SlowControlC return false; } + // fewer, larger packets are better for network performance, so we batch logging and monitoring messages. + // on the other hand, if packet size exceeds the MTU, they will fragment, increasing packets and reducing reliability + // so, try to batch up to the MTU size. In order to do that, we need to know the MTU. + //std::set interfaces = GetInterfaces(); + //if(interfaces.size()) MSS_SIZE = GetMTU(interfaces[??]); // but which interface? + // just get MTU from config variable. *sigh* + return true; } @@ -978,24 +987,35 @@ bool Services::SendLog(const std::string& message, LogLevel severity, const std: const std::string& name = (device=="") ? m_name : device; - // FIXME we should be able to relax this check if compression is enabled... - if((message.length()+name.length())>MAX_MSG_SIZE){ - if(m_verbose) std::cerr<<"Logging message is too long!"< locker(logging_buf_mtx); + // merge identical logging messages back-to-back if(logging_buf.size() && name==logging_buf.back().device && message==logging_buf.back().message){ ++logging_buf.back().repeats; return true; } - // grab timestamp at time of call if 0 - time_t ts = (timestamp!=0) ? timestamp : time(nullptr)*1000; + // reject if this message is too big to fit in a UDP datagram even with compression + size_t compressed_bytes = ZSTD_compressBound(message.length()+name.length()); + if(compressed_bytes > MAX_MSG_SIZE){ + if(m_verbose) std::cerr<<"Logging message is too long!"< MSS_SIZE){ + BatchAndSendMulticast(&thread_args, false, true); // don't try to lock logging buffer, we've got it + } logging_buf.emplace_back(message, severity, name, ts); + thread_args.logging_batch_bytes += compressed_bytes; + return true; } @@ -1017,24 +1037,38 @@ bool Services::SendMonitoringData(const std::string& json_data, const std::strin const std::string& name = (device=="") ? m_name : device; - if((json_data.length()+name.length()+subject.length())>MAX_MSG_SIZE){ - if(m_verbose) std::cerr<<"Monitoring message is too long!"< locker(monitoring_buf_mtx); - // take first of repeated monitoring sends within buffer period + // only accept the first of repeated monitoring sends within buffer period auto it = monitoring_buf.find(name+subject); if(it!=monitoring_buf.end() && (ts - it->second.timestamp) MAX_MSG_SIZE){ + if(m_verbose) std::cerr<<"Monitoring message is too long!"< MSS_SIZE){ + BatchAndSendMulticast(&thread_args, true, false); // don't try to lock monitoring buffer, we've got it + } + monitoring_buf.emplace(std::piecewise_construct, std::forward_as_tuple(name+subject), std::forward_as_tuple(json_data, subject, name, ts)); + thread_args.monitoring_batch_bytes += compressed_bytes; + return true; } @@ -1062,8 +1096,8 @@ bool Services::SendROOTplotMulticast(const std::string& plot_name, const std::st + ", \"lifetime\":"+std::to_string(lifetime) + ", \"data\":"+ json_data+"}"; - if(cmd_string.length()>MAX_UDP_PACKET_SIZE){ - if(m_verbose) std::cerr<<"ROOT plot json is too long! Maximum length may be MAX_UDP_PACKET_SIZE bytes"< MAX_UDP_PACKET_SIZE){ + if(m_verbose) std::cerr<<"ROOT plot json is too long!"<(args); + BatchAndSendMulticast(m_args, true, true); + std::this_thread::sleep_until(m_args->last_send+m_args->multicast_send_period_ms); + return; + +} + +bool Services::BatchAndSendMulticast(BufferThreadArgs* m_args, bool log_lock, bool mon_lock){ m_args->last_send = std::chrono::steady_clock::now(); m_args->local_merge_buf.clear(); - std::unique_lock locker(*m_args->logging_buf_mtx); + std::unique_lock locker(*m_args->logging_buf_mtx, std::defer_lock); + if(log_lock) locker.lock(); // merge into a batch bool first=true; @@ -1274,13 +1316,16 @@ void Services::BufferThread(Thread_args* args){ } // send - if(m_args->local_merge_buf.empty() || m_args->services->SendLog(m_args->local_merge_buf)){ - m_args->logging_buf->clear(); // FIXME do we not clear on error...? does it depend on the error...? + if(!m_args->local_merge_buf.empty()){ + m_args->services->SendLog(m_args->local_merge_buf); + m_args->logging_buf->clear(); + m_args->logging_batch_bytes = 0; } // repeat for monitoring messages m_args->local_merge_buf.clear(); - locker = std::unique_lock(*m_args->monitoring_buf_mtx); + locker = std::unique_lock(*m_args->monitoring_buf_mtx, std::defer_lock); + if(mon_lock) locker.lock(); first=true; for(std::pair& msg : *m_args->monitoring_buf){ @@ -1295,8 +1340,10 @@ void Services::BufferThread(Thread_args* args){ } // send - if(m_args->local_merge_buf.empty() || m_args->services->SendMonitoringData(m_args->local_merge_buf)){ - m_args->monitoring_buf->clear(); // FIXME do we not clear on error...? does it depend on the error...? + if(!m_args->local_merge_buf.empty()){ + m_args->services->SendMonitoringData(m_args->local_merge_buf); + m_args->monitoring_buf->clear(); + m_args->monitoring_batch_bytes = 0; } // our other sevice task: prune the alarm buffer. @@ -1309,12 +1356,7 @@ void Services::BufferThread(Thread_args* args){ else ++it; } - // release mtx - locker.unlock(); - - std::this_thread::sleep_until(m_args->last_send+m_args->multicast_send_period_ms); - - return; + return true; } std::string Services::JsonEscape(std::string s){ @@ -1335,7 +1377,7 @@ std::string Services::GetLocalConfig(){ std::string Services::SCLocalConfig(const char* data){ - return "["+std::to_string(m_base_config_id)+","+std::to_string(m_run_mode_config_id)+"]: "+m_local_config; + return "base: "+std::to_string(m_base_config_id)+", runmode:"+std::to_string(m_run_mode_config_id)+", config: "+m_local_config; } @@ -1380,18 +1422,18 @@ bool Services::SetChangeConfigFunc(std::function func){ (*sc_vars)["Config"]->SetValue((int)ConfigState::ChangeStart); bool ok = func(m_local_config); if(!ok){ - std::cerr<<"ChangeConfig Error"<SetWarning(true); - } + std::cerr<<"ChangeConfig Error"<SetWarning(true); + } int new_state = ok ? (int)ConfigState::ChangeEnd : (int)ConfigState::ChangeFail; (*sc_vars)["Config"]->SetValue(new_state); (*sc_vars)["NewConfig"]->SetValue(0); return (ok ? "OK" : "Error"); - }, - 0, - false); // this version will not be locked during non-testing runs, - // since it only allows loading configurations in line with the current run type. + }, // setter + 0, // getter + false, // not locked during non-testing runs, as it only allows loading configurations in line with the current run type + false); // not hidden // 3. allgood = allgood && @@ -1403,15 +1445,15 @@ bool Services::SetChangeConfigFunc(std::function func){ int new_state = ok ? (int)ConfigState::ChangeEnd : (int)ConfigState::ChangeFail; (*sc_vars)["Config"]->SetValue(new_state); if(!ok){ - sc_vars->SetWarning(true); - std::cerr<<"ChangeConfig Error"<SetWarning(true); + std::cerr<<"ChangeConfig Error"<