From a89871d73ebbff0349e6728cbc9d37e33346d35f Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Tue, 12 May 2026 14:24:31 +0000 Subject: [PATCH] replace stringstream zmq parsing with string construction so we don't rely on null terminated messages --- UserTools/Logger/Logger.cpp | 6 ++-- UserTools/template/MyToolZMQMultiThread.cpp | 6 ++-- src/DAQDataModelBase/DAQUtilities.h | 3 +- src/DAQLogging/DAQLogging.cpp | 27 ++++++++--------- src/NodeDaemon/NodeDaemon.cpp | 13 +++++---- src/RemoteControl/MCDebug.cpp | 4 +-- src/RemoteControl/RemoteControl.cpp | 17 ++++------- src/ServiceDiscovery/ServiceDiscovery.cpp | 17 ++++++----- .../SlowControlCollection.cpp | 29 +++++++++---------- src/ToolDAQChain/ToolDAQChain.cpp | 3 +- 10 files changed, 62 insertions(+), 63 deletions(-) diff --git a/UserTools/Logger/Logger.cpp b/UserTools/Logger/Logger.cpp index dbcee5ca..37b30271 100644 --- a/UserTools/Logger/Logger.cpp +++ b/UserTools/Logger/Logger.cpp @@ -46,13 +46,13 @@ bool Logger::Execute(){ zmq::message_t Rmessage; if( LogReceiver->recv (&Rmessage)){ // printf("got a message \n"); - std::istringstream ss(static_cast(Rmessage.data())); + std::string ss(static_cast(Rmessage.data()),Rmessage.size()); - *m_log<recv(&message); - std::istringstream iss(static_cast(message.data())); - *m_log<<"reply = "<(message.data()),message.size()); + *m_log<<"reply = "<ThreadReceive->recv(&message); - std::istringstream iss(static_cast(message.data())); + std::string ss(static_cast(message.data()),message.size()); sleep(10); diff --git a/src/DAQDataModelBase/DAQUtilities.h b/src/DAQDataModelBase/DAQUtilities.h index cfd66735..ca42ec97 100644 --- a/src/DAQDataModelBase/DAQUtilities.h +++ b/src/DAQDataModelBase/DAQUtilities.h @@ -119,7 +119,8 @@ namespace ToolFramework{ if(sock->recv(&message)){ - std::istringstream iss(static_cast(message.data())); + std::string ss(static_cast(message.data()),message.size()); + std::istringstream iss(ss); // long long unsigned int tmpP; unsigned long tmpP; diff --git a/src/DAQLogging/DAQLogging.cpp b/src/DAQLogging/DAQLogging.cpp index 35f68079..27d9de44 100644 --- a/src/DAQLogging/DAQLogging.cpp +++ b/src/DAQLogging/DAQLogging.cpp @@ -287,16 +287,16 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw zmq::message_t Receive; LogReceiver.recv (&Receive); - std::istringstream ss(static_cast(Receive.data())); + std::string ss(static_cast(Receive.data()),Receive.size()); if (logfile.is_open()) { - logfile << ss.str();//<(Receive.data())); + std::string ss(static_cast(Receive.data()),Receive.size()); - if(ss.str()=="Quit"){ + if(ss=="Quit"){ //printf("%s \n","received quit"); running=false; } @@ -388,13 +388,13 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw outmessage.Set("msg_id",msg_id); outmessage.Set("msg_time", isot.str()); outmessage.Set("msg_type", "Log"); - outmessage.Set("msg_value",ss.str()); + outmessage.Set("msg_value",ss); */ outmessage.Set("topic","logging"); outmessage.Set("time", isot.str()); outmessage.Set("device",args->m_service); outmessage.Set("severity","logging"); - outmessage.Set("message",ss.str()); + outmessage.Set("message",ss); std::string rmessage; outmessage>>rmessage; @@ -459,7 +459,7 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw zmq::message_t Receive; if(LogReceiver.recv (&Receive)){ - std::istringstream ss(static_cast(Receive.data())); + std::string ss(static_cast(Receive.data()),Receive.size()); boost::posix_time::ptime t = boost::posix_time::microsec_clock::universal_time(); @@ -474,7 +474,7 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw outmessage.Set("msg_id",msg_id); *outmessage["msg_time"]=isot.str(); *outmessage["msg_type"]="Log"; - outmessage.Set("msg_value",ss.str()); + outmessage.Set("msg_value",ss); for(std::map::iterator it=RemoteConnections.begin(); it!=RemoteConnections.end(); ++it){ @@ -500,7 +500,7 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw } - if(ss.str()=="Quit"){ + if(ss=="Quit"){ //printf("%s \n","received quit"); running=false; } @@ -536,7 +536,8 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw //printf("sent sd req \n"); zmq::message_t receive; if(Ireceive.recv(&receive)){ - std::istringstream iss(static_cast(receive.data())); + std::string ss(static_cast(receive.data()),receive.size()); + std::istringstream iss(ss); //printf("received from sd \n"); int size; @@ -560,8 +561,8 @@ src/DAQLogging/DAQLogging.{h,cpp} -nw zmq::message_t servicem; Ireceive.recv(&servicem); - std::istringstream ss(static_cast(servicem.data())); - service->JsonParser(ss.str()); + std::string ss2(static_cast(servicem.data()),servicem.size()); + service->JsonParser(ss2); std::string servicetype; std::string uuid; diff --git a/src/NodeDaemon/NodeDaemon.cpp b/src/NodeDaemon/NodeDaemon.cpp index 72d10658..ce04cd49 100644 --- a/src/NodeDaemon/NodeDaemon.cpp +++ b/src/NodeDaemon/NodeDaemon.cpp @@ -140,11 +140,11 @@ int main(int argc, char* argv[]){ direct.recv(&message); - std::istringstream iss(static_cast(message.data())); - //std::cout<<"Received message: "<(message.data()), message.size()); + //std::cout<<"Received message: "<("var1").c_str()); + std::ofstream outfile (bb.Get("var1").c_str(), std::ios::binary); if (outfile.is_open()){ while (1) { @@ -207,14 +207,15 @@ int main(int argc, char* argv[]){ int64_t more; size_t size = sizeof(int64_t); ftp.getsockopt(ZMQ_RCVMORE, &more, &size); - std::istringstream filess(static_cast(file.data())); + //std::istringstream filess(static_cast(file.data()); // char *tmp=static_cast(file.data()); // std::cout<<"buf before = "<>tmp; //std::string tmp2; // filess>>tmp2; //std::cout<<"received = "<(file.data()),file.size()); + outfile<<"\n"; // std::cout<<"received part"< 0){ - printf("%s: message = \"%s\"\n", inet_ntoa(addr.sin_addr), message); + printf("%s: message = \"%.*s\"\n", inet_ntoa(addr.sin_addr), cnt, message); //if(message[0]!='[') break; @@ -112,4 +112,4 @@ int main(){ } } - \ No newline at end of file + diff --git a/src/RemoteControl/RemoteControl.cpp b/src/RemoteControl/RemoteControl.cpp index 30698a3e..849f6293 100644 --- a/src/RemoteControl/RemoteControl.cpp +++ b/src/RemoteControl/RemoteControl.cpp @@ -78,7 +78,8 @@ int main(int argc, char** argv){ zmq::message_t receive; Ireceive.recv(&receive); - std::istringstream iss(static_cast(receive.data())); + std::string ss(static_cast(receive.data()),receive.size()); + std::istringstream iss(ss); int size; iss>>size; @@ -97,8 +98,8 @@ int main(int argc, char** argv){ zmq::message_t servicem; Ireceive.recv(&servicem); - std::istringstream ss(static_cast(servicem.data())); - service->JsonParser(ss.str()); + std::string ss(static_cast(servicem.data()),servicem.size()); + service->JsonParser(ss); std::string name; name=(*service).Get("msg_value"); @@ -236,10 +237,7 @@ 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(); + std::string answer(static_cast(receive.data()),receive.size()); Store rr; rr.JsonParser(answer); @@ -364,10 +362,7 @@ 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(); + std::string answer(static_cast(receive.data()),receive.size()); Store rr; rr.JsonParser(answer); diff --git a/src/ServiceDiscovery/ServiceDiscovery.cpp b/src/ServiceDiscovery/ServiceDiscovery.cpp index 1150d16a..652d0c83 100644 --- a/src/ServiceDiscovery/ServiceDiscovery.cpp +++ b/src/ServiceDiscovery/ServiceDiscovery.cpp @@ -169,7 +169,8 @@ void* ServiceDiscovery::MulticastPublishThread(void* arg){ zmq::message_t commands; Ireceive.recv(&commands); - std::istringstream tmp(static_cast(commands.data())); + std::string stmp(static_cast(commands.data()),commands.size()); + std::istringstream tmp(stmp); std::string command; std::string service; boost::uuids::uuid uuid; @@ -324,9 +325,9 @@ void* ServiceDiscovery::MulticastPublishThread(void* arg){ if(in[0].revents & ZMQ_POLLIN){ zmq::message_t Ereceive; StatusCheck.recv (&Ereceive); - std::istringstream ss(static_cast(Ereceive.data())); + std::string ss(static_cast(Ereceive.data()),Ereceive.size()); - mm.JsonParser(ss.str()); + mm.JsonParser(ss); } } } @@ -580,7 +581,7 @@ void* ServiceDiscovery::MulticastListenThread(void* arg){ Store* newservice= new Store(); newservice->Set("ip",inet_ntoa(addr.at(i).sin_addr)); - newservice->JsonParser(message); + newservice->JsonParser(std::string(&message[0],cnt)); std::string uuid; newservice->Get("uuid",uuid); @@ -665,7 +666,8 @@ void* ServiceDiscovery::MulticastListenThread(void* arg){ if(Ireceive.recv(&comm)){ - std::istringstream iss(static_cast(comm.data())); + std::string ss(static_cast(comm.data()),comm.size()); + std::istringstream iss(ss); std::string arg1=""; std::string arg2=""; @@ -676,9 +678,10 @@ void* ServiceDiscovery::MulticastListenThread(void* arg){ //printf("d2\n"); //zmq::message_t sizem(512); int size= RemoteServices.size(); - zmq::message_t sizem(sizeof size); + std::string sizes = std::to_string(size); + zmq::message_t sizem(sizes.length()+1); - snprintf ((char *) sizem.data(), sizeof size , "%d" ,size) ; + snprintf ((char *) sizem.data(), sizes.length()+1 , "%d" ,size) ; // zmq::poll(out,1,1000); diff --git a/src/ServiceDiscovery/SlowControlCollection.cpp b/src/ServiceDiscovery/SlowControlCollection.cpp index 56e1c868..ed8eb3c5 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -244,10 +244,10 @@ void SlowControlCollection::Thread(Thread_args* arg){ return; } - std::istringstream iss(static_cast(message.data())); + std::string iss(static_cast(message.data()),message.size()); Store tmp; - //printf("iss=%s\n",iss.str().c_str()); - tmp.JsonParser(iss.str()); + //printf("iss=%s\n",iss.c_str()); + tmp.JsonParser(iss); //tmp.Print(); if(!tmp.Has("msg_value")){ std::cerr<<"error: Poorly formatted slowcontrol input [no msg_value]"<(message.data())); + std::string iss(static_cast(message.data()),message.size()); // receive alert payload std::string payload; @@ -363,15 +363,14 @@ void SlowControlCollection::Thread(Thread_args* arg){ // FIXME this case should be handled! what do we do? std::cerr<<"failed to receive alert payload!"<(message.data()),message.size()); has_data=true; } - //std::cout<alert_functions_mutex->lock(); - if(iss.str() == "LoadConfig") (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadStart); - else if(iss.str() == "ChangeConfig"){ + if(iss == "LoadConfig") (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadStart); + else if(iss == "ChangeConfig"){ if((*args->SC_vars)["NewConfig"]->GetValue() == 0) return; (*args->SC_vars)["Config"]->SetValue((int)ConfigState::ChangeStart); } @@ -379,19 +378,19 @@ void SlowControlCollection::Thread(Thread_args* arg){ bool error = false; - if(args->alert_functions->count(iss.str())){ + if(args->alert_functions->count(iss)){ if(has_data){ try{ - error = !((*(args->alert_functions))[iss.str()](iss.str().c_str(), payload.c_str())); + error = !((*(args->alert_functions))[iss](iss.c_str(), payload.c_str())); } catch(...){ error = true; } - if(iss.str() == "LoadConfig"){ + if(iss == "LoadConfig"){ if(error)(*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadFail); else (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadEnd); } - else if(iss.str() == "ChangeConfig"){ + else if(iss == "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); @@ -400,13 +399,13 @@ void SlowControlCollection::Thread(Thread_args* arg){ } else try{ - error=!((*(args->alert_functions))[iss.str()](iss.str().c_str(), 0)); + error=!((*(args->alert_functions))[iss](iss.c_str(), 0)); } catch(...){ error = true; } - if(error) std::cerr<<"alert fucntion failed: "<(message.data())); - command=iss.str(); + command = std::string(static_cast(message.data()),message.size()); Store rr; rr.JsonParser(command);