From dac67adc714c1cab7430bdb30bc52280093ea32d Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Mon, 10 Aug 2026 12:19:47 +0100 Subject: [PATCH 1/2] 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 f339179..ea06a32 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) } @@ -620,16 +630,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++; + } + } } @@ -641,17 +651,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--; } } } @@ -685,8 +695,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; @@ -695,7 +705,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"); @@ -704,49 +714,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; @@ -848,3 +858,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 2192f24c8aa63aefdfcb2304045045e5a598f4fa Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Mon, 10 Aug 2026 12:45:09 +0100 Subject: [PATCH 2/2] 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 ea06a32..0eca480 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -886,7 +886,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);