From 35f8dd622da4c07c1a828e96a63c183af831b156 Mon Sep 17 00:00:00 2001 From: Marcus O'Flaherty Date: Thu, 10 Sep 2026 17:59:22 +0100 Subject: [PATCH] * make wait on startup for services to connect to middleman configurable * forward Services class to SlowControlCollection so it can log problems * add a separate buffer of monitoring (device+subject) and reject duplicates in monitoring merge period based on this, rather than the buffer. Needed as the buffer may be flushed early if we exceed the MTU size, which would otherwise bypass rate limiting. * set "Config" slow control state in LoadConfigSlowControlFunc * add conversion of line breaks in logging messages to '\n', as these otherwise break DB insertion! * report state of testing flag in LocalConfig slow control response * set SlowControlCollection warning bit if a slow control callback fails... tenatively * set SlowControlCollection error bit if an alert fails to receive, decompress, or callback function fails * make Services a friend class of SlowControlCollection, and make SetTesting, TestingEnable, TestingDisable functions private --- src/ServiceDiscovery/Services.cpp | 32 +++-- src/ServiceDiscovery/Services.h | 4 +- src/ServiceDiscovery/ServicesBackend.cpp | 2 +- .../SlowControlCollection.cpp | 129 ++++++++++++------ src/ServiceDiscovery/SlowControlCollection.h | 17 ++- 5 files changed, 122 insertions(+), 62 deletions(-) diff --git a/src/ServiceDiscovery/Services.cpp b/src/ServiceDiscovery/Services.cpp index d93948d4..de163fca 100644 --- a/src/ServiceDiscovery/Services.cpp +++ b/src/ServiceDiscovery/Services.cpp @@ -16,7 +16,7 @@ Services::Services(){ m_run_mode_config_id=0; } - + Services::~Services(){ // kill background buffering thread @@ -36,6 +36,7 @@ bool Services::Init(Store &m_variables, zmq::context_t* context_in, SlowControlC m_context = context_in; sc_vars = sc_vars_in; + sc_vars->SetServices(this); bool alerts_send = 0; int alert_send_port = 12242; @@ -88,7 +89,11 @@ bool Services::Init(Store &m_variables, zmq::context_t* context_in, SlowControlC // so we need to wait for the middleman to receive one & connect before we can communicate with it. int pub_period=5; m_variables.Get("service_publish_sec",pub_period); - if(!Ready(pub_period*3000)){ // Wait up to 3 broadcast periods. It'll return sooner if it connects. + // Wait up to 3 broadcast periods by default. It'll return sooner if it connects. + int ready_wait_ms = pub_period*3000; + m_variables.Get("ready_wait_ms",ready_wait_ms); + + if(ready_wait_ms>0 && !Ready(ready_wait_ms)){ if(m_verbose) std::cerr<<"Warning: service not yet connected..."<& resp std::string err=""; - // for now, commands must not begin with '(' as this is used to identify compressed messages. - // Since we don't know what a user-provided query string may be, prepend with a space to ensure this. - std::string sanitized_query = std::string{" "}+query; - - if(!m_backend_client.SendCommand("W_QUERY", sanitized_query, &responses, timeout, &err)){ + if(!m_backend_client.SendCommand("W_QUERY", query, &responses, timeout, &err)){ if(m_verbose) std::cerr<<"SQLQuery error: "< locker(monitoring_buf_mtx); // 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)second)SetValue((int)ConfigState::LoadStart); bool success = LoadConfigAlertFunc("",payload); + if(success)(*sc_vars)["Config"]->SetValue((int)ConfigState::LoadEnd); + else (*sc_vars)["Config"]->SetValue((int)ConfigState::LoadFail); if(!success) return std::string("Failed to load config: ")+payload; return std::string("Loaded config: ")+payload; @@ -1318,6 +1327,7 @@ bool Services::BatchAndSendMulticast(BufferThreadArgs* m_args, bool log_lock, bo m_args->services->SendMonitoringData(m_args->local_merge_buf); m_args->monitoring_buf->clear(); m_args->monitoring_batch_bytes = 0; + if(mon_lock) m_args->monitoring_msgs_sent->clear(); } // our other sevice task: prune the alarm buffer. @@ -1336,8 +1346,10 @@ bool Services::BatchAndSendMulticast(BufferThreadArgs* m_args, bool log_lock, bo std::string Services::JsonEscape(std::string s){ // TODO is there a more efficient way to do this... std::string out; + while(s.size() && std::isspace(s.back())) s.pop_back(); for(char& a : s){ if(a=='"' || a=='\\') out.push_back('\\'); + if(a=='\x0a' || a=='\x0d'){ out+="\\n"; continue; } out.push_back(a); } return out; @@ -1351,7 +1363,7 @@ std::string Services::GetLocalConfig(){ std::string Services::SCLocalConfig(const char*){ - return "base: "+std::to_string(m_base_config_id)+", runmode:"+std::to_string(m_run_mode_config_id)+", config: "+m_local_config; + return "base: "+std::to_string(m_base_config_id)+", runmode:"+std::to_string(m_run_mode_config_id)+", testing:"+std::to_string(m_testing)+", config: "+m_local_config; } diff --git a/src/ServiceDiscovery/Services.h b/src/ServiceDiscovery/Services.h index 5ae01771..51a8e1f4 100644 --- a/src/ServiceDiscovery/Services.h +++ b/src/ServiceDiscovery/Services.h @@ -56,6 +56,7 @@ namespace ToolFramework { Services* services; std::vector* logging_buf; std::unordered_map* monitoring_buf; + std::unordered_map* monitoring_msgs_sent; std::vector* alarm_buf; std::mutex* logging_buf_mtx; std::mutex* monitoring_buf_mtx; @@ -164,11 +165,12 @@ namespace ToolFramework { std::vector logging_buf; std::unordered_map monitoring_buf; + std::unordered_map monitoring_msgs_sent; std::vector alarm_buf; std::mutex logging_buf_mtx; std::mutex monitoring_buf_mtx; std::mutex alarm_buf_mtx; - uint32_t mon_merge_period_ms; + uint32_t mon_merge_period_ms; // N.B. can be at most multicast_send_period_ms uint32_t multicast_send_period_ms; uint32_t alarm_cooldown_ms; diff --git a/src/ServiceDiscovery/ServicesBackend.cpp b/src/ServiceDiscovery/ServicesBackend.cpp index bbaf6fc3..bd1d925a 100644 --- a/src/ServiceDiscovery/ServicesBackend.cpp +++ b/src/ServiceDiscovery/ServicesBackend.cpp @@ -727,7 +727,7 @@ bool ServicesBackend::GetNextResponse(){ cmd.success = false; Log(cmd.err, v_warning, m_verbosity); break; - } else { + } else if(cmd.success){ cmd.response.push_back(std::move(cmd.err)); cmd.err = ""; } diff --git a/src/ServiceDiscovery/SlowControlCollection.cpp b/src/ServiceDiscovery/SlowControlCollection.cpp index f44a9aed..4c541313 100644 --- a/src/ServiceDiscovery/SlowControlCollection.cpp +++ b/src/ServiceDiscovery/SlowControlCollection.cpp @@ -1,5 +1,6 @@ #include #include "zstd_helpers.h" +#include "Services.h" using namespace ToolFramework; @@ -12,6 +13,7 @@ SlowControlCollectionThread_args::SlowControlCollectionThread_args(){ alert_functions=0; alert_functions_mutex=0; SC_vars=0; + m_services=0; } @@ -28,6 +30,7 @@ SlowControlCollectionThread_args::~SlowControlCollectionThread_args(){ alert_functions=0; alert_functions_mutex=0; SC_vars=0; + m_services=0; if(pub_monitor_socket){ if(m_pub) zmq_socket_monitor((void*)(*m_pub), NULL, 0); // stop pub socket sending events @@ -43,6 +46,7 @@ SlowControlCollectionThread_args::~SlowControlCollectionThread_args(){ SlowControlCollection::SlowControlCollection(){ args=0; + m_services=0; m_util=0; m_context=0; m_pub=0; @@ -58,6 +62,10 @@ SlowControlCollection::SlowControlCollection(){ } +void SlowControlCollection::SetServices(Services* services){ + m_services = services; +} + SlowControlCollection::~SlowControlCollection(){ Stop(); @@ -93,6 +101,7 @@ void SlowControlCollection::Stop(){ //printf("p6\n"); Clear(); //printf("p7\n"); + m_services=0; } bool SlowControlCollection::Init(zmq::context_t* context, int sc_port, bool new_service, int alert_receive_port, bool alerts_receive, int alert_send_port, bool alerts_send){ @@ -131,14 +140,16 @@ bool SlowControlCollection::Init(zmq::context_t* context, int sc_port, bool new_ delete args; args=0; - std::clog<<"Error adding alert send port to SD"<SendLog("Error adding alert send port to SD", LogLevel::Error); + SetError(true); return false; } if(m_thread){ // tell the socket to send monitoring messages so we can identify when its connected if(zmq_socket_monitor((void*)(*m_pub), "inproc://AlertSendMonitor", ZMQ_EVENT_ALL)!=0){ - std::cerr<<"SCC error starting monitor on alert send socket: "<SendLog(std::string{"SCC error starting monitor on alert send socket: "}+zmq_strerror(errno), LogLevel::Error); + SetError(true); } // open a pair socket to receive them pub_monitor_socket = new zmq::socket_t(*context, ZMQ_PAIR); @@ -174,7 +185,8 @@ bool SlowControlCollection::Init(zmq::context_t* context, int sc_port, bool new_ args=0; - std::clog<<"Error adding port alert receive to SD"<SendLog("Error adding port alert receive to SD",LogLevel::Error); + SetError(true); return false; } @@ -197,7 +209,8 @@ bool SlowControlCollection::Init(zmq::context_t* context, int sc_port, bool new_ delete args; args=0; - std::clog<<"Error adding port SC to SD"<SendLog("Error adding port SC to SD",LogLevel::Error); + SetError(true); return false; } @@ -270,7 +283,7 @@ void SlowControlCollection::Thread(Thread_args* arg){ int ok = args->sock->recv(&identity); if(ok==0 || !identity.more()){ - std::cerr<<"error: Poorly formatted slowcontrol input [identity problem]"<m_services->SendLog("Poorly formatted slowcontrol input [identity problem]",LogLevel::Warning); return; } @@ -278,7 +291,7 @@ void SlowControlCollection::Thread(Thread_args* arg){ ok = args->sock->recv(&blank); if (!blank.more()){ - std::cerr<<"error: Poorly formatted slowcontrol input [blank problem]"<m_services->SendLog("Poorly formatted slowcontrol input [blank problem]",LogLevel::Warning); return; } @@ -286,14 +299,14 @@ void SlowControlCollection::Thread(Thread_args* arg){ ok = args->sock->recv(&message); if(ok==0 || message.more()){ - std::cerr<<"error: Poorly formatted slowcontrol input [message problem]"<m_services->SendLog("Poorly formatted slowcontrol input [message problem]",LogLevel::Warning); return; } std::string payload; std::unique_lock locker(args->SCC->zstd_dctx_mtx); if(!ZstdDecompress(args->SCC->zstd_dctx, (char*)message.data(), message.size(), payload, args->SCC->MAX_DECOMPRESSED_SIZE)){ - std::cerr<<"failed to decompress slow control message: "<m_services->SendLog("failed to decompress slow control message",LogLevel::Warning); return; } locker.unlock(); @@ -302,7 +315,7 @@ void SlowControlCollection::Thread(Thread_args* arg){ tmp.JsonParser(payload); //tmp.Print(); if(!tmp.Has("msg_value")){ - std::cerr<<"error: Poorly formatted slowcontrol input [no msg_value]"<m_services->SendLog("Poorly formatted slowcontrol input [no msg_value]",LogLevel::Warning); return; } @@ -322,7 +335,12 @@ void SlowControlCollection::Thread(Thread_args* arg){ std::string reply=""; bool strip=false; - Update(args->SCC, key, value, reply, strip, *(args->testing)); + bool callback_successs = Update(args->SCC, key, value, reply, strip, *(args->testing)); + if(!callback_successs){ + args->SCC->SetWarning(true); + // is this generally appropriate? Should we SetError or SetWarning? should we leave it to users? + // users may have called SetError within their callback...this would override that... + } /* if(key == "?"){ @@ -389,11 +407,9 @@ void SlowControlCollection::Thread(Thread_args* arg){ if(tmp_ok) tmp_ok = tmp_ok && args->sock->send(blank, ZMQ_SNDMORE); if(tmp_ok) tmp_ok= tmp_ok && args->sock->send(zmsg); if(!tmp_ok){ - std::cerr<<"failed to send '"<m_services->SendLog("failed to send SlowControl '"+key+"' reply '"+reply+"'",LogLevel::Warning); return; } - // FIXME these sorts of errors should be logged somewhere - // rather than being silently ignored. This info could be critical for debugging issues!!! } if (args->alerts_receive && args->items[1].revents & ZMQ_POLLIN){ //received alert value; @@ -403,8 +419,9 @@ void SlowControlCollection::Thread(Thread_args* arg){ // receive alert type int ok = args->sub->recv(&message); if(ok==0){ - // FIXME this case should be handled! what do we do? - std::cerr<<"failed to receive alert!"<m_services->SendLog("error receiving alert!",LogLevel::Error); + args->SCC->SetError(true); + return; } std::istringstream iss(static_cast(message.data())); @@ -414,23 +431,22 @@ void SlowControlCollection::Thread(Thread_args* arg){ if(message.more()){ ok = args->sub->recv(&message); if(ok==0){ - // FIXME this case should be handled! what do we do? - std::cerr<<"failed to receive "<m_services->SendLog("failed to receive "+iss.str()+"' alert payload",LogLevel::Warning); + args->SCC->SetError(true); return; } if(!ZstdDecompress(args->SCC->zstd_dctx, (char*)message.data(), message.size(), payload)){ - std::cerr<<"failed to decompress "<m_services->SendLog("failed to decompress "+iss.str()+"' alert payload",LogLevel::Warning); return; } has_data=true; } - //int a=0; + int a=0; while(message.more()){ - - args->sub->recv(&message); // FIXME do we want any warnings or handling here? - //memcpy((void*)payload.data(),message.data(),message.size()); - //a++; + args->sub->recv(&message); + if(a==0) args->m_services->SendLog("Unexpected additional "+iss.str()+" alert parts",LogLevel::Warning); + a++; } if(iss.str() == "LoadConfig") (*args->SC_vars)["Config"]->SetValue((int)ConfigState::LoadStart); @@ -473,28 +489,34 @@ void SlowControlCollection::Thread(Thread_args* arg){ } } - if(error) std::cerr<<"alert function failed: "<m_services->SendLog("alert function failed: "+iss.str(),LogLevel::Error); + args->SCC->SetError(true); + } } if(args->alert_functions->count("*")){ if(has_data){ - try{ - error = !((*(args->alert_functions))["*"](iss.str().c_str(), payload.c_str())); - } - catch(...){ - error = true; - } + try{ + error = !((*(args->alert_functions))["*"](iss.str().c_str(), payload.c_str())); + } + catch(...){ + error = true; + } } else { - try{ - error=!((*(args->alert_functions))["*"](iss.str().c_str(), 0)); - } - catch(...){ - error = true; - } + try{ + error=!((*(args->alert_functions))["*"](iss.str().c_str(), 0)); + } + catch(...){ + error = true; + } + } + + if(error){ + args->m_services->SendLog("alert function failed: "+iss.str(),LogLevel::Error); + args->SCC->SetError(true); } - - if(error) std::cerr<<"alert function failed: "<SendLog("error sending alert: functionality not enabled",LogLevel::Error); + return false; + } + bool ok; + zmq::message_t message(alert.length()+1); snprintf((char*) message.data(), alert.length()+1, "%s", alert.c_str()); if(payload==""){ - return m_pub->send(message); // err: "zmq send "+zmq_strerror(errno) + ok = m_pub->send(message); + if(!ok){ + m_services->SendLog("zmq error sending alert '"+alert+"': "+zmq_strerror(errno),LogLevel::Error); + return false; + } + return true; } + // if we didn't return, we have a payload as well - bool ok = m_pub->send(message, ZMQ_SNDMORE); - if(!ok) return false; // err: "zmq send "+zmq_strerror(errno) + ok = m_pub->send(message, ZMQ_SNDMORE); + if(!ok){ + m_services->SendLog("zmq error sending alert '"+alert+"': "+zmq_strerror(errno),LogLevel::Error); + return false; + } std::unique_locklocker(args->SCC->zstd_cctx_mtx); std::string compress_buffer; std::pair output = ZstdCompress(args->SCC->zstd_cctx, payload.data(), payload.size(), compress_buffer); zmq::message_t message2(output.second); memcpy(message2.data(), output.first, output.second); - return m_pub->send(message2); // err: "zmq send "+zmq_strerror(errno) - + ok = m_pub->send(message2); + if(!ok){ + m_services->SendLog("zmq error sending alert '"+alert+"': "+zmq_strerror(errno),LogLevel::Error); + return false; + } + return true; } void SlowControlCollection::JsonParser(std::string json){ diff --git a/src/ServiceDiscovery/SlowControlCollection.h b/src/ServiceDiscovery/SlowControlCollection.h index be7411e1..d2766821 100644 --- a/src/ServiceDiscovery/SlowControlCollection.h +++ b/src/ServiceDiscovery/SlowControlCollection.h @@ -9,7 +9,9 @@ #include namespace ToolFramework{ - + + class Services; + //typedef void (*AlertFunction)(const char*, const char*); typedef std::function AlertFunction; @@ -35,6 +37,7 @@ namespace ToolFramework{ bool alerts_receive; bool* testing; std::map* SC_vars; + Services* m_services = nullptr; // for checking when alert send socket is ready zmq::socket_t* m_pub = nullptr; @@ -45,6 +48,9 @@ namespace ToolFramework{ class SlowControlCollection{ + // allow Services class to access e.g. SetTesting function + friend class Services; + public: SlowControlCollection(); @@ -54,6 +60,7 @@ namespace ToolFramework{ bool ListenForData(int poll_length=0); bool InitThreadedReceiver(zmq::context_t* context, int port=60000, int poll_length=100, bool new_service=true, int alert_receive_port=12243, bool alert_receive=true, int alert_send_port=12242, bool alert_send=true); SlowControlElement* operator[](std::string key); + void SetServices(Services*); bool Add(std::string name, SlowControlElementType type, SCFunction change_function = 0, SCFunction read_function = 0, bool testing_lock=true, bool hidded=false); bool Remove(std::string name); void Clear(); @@ -63,10 +70,7 @@ namespace ToolFramework{ std::string PrintJSON(); void Stop(); void JsonParser(std::string json); - void SetTesting(bool testing); bool GetTesting(); - void TestingEnable(); - void TestingDisable(); bool Ready(int timeout_ms); void SetActive(bool active); void SetError(bool error); @@ -78,8 +82,10 @@ namespace ToolFramework{ return SC_vars[name]->GetValue(); } - private: + void SetTesting(bool testing); + void TestingEnable(); + void TestingDisable(); std::map SC_vars; std::map m_alert_functions; @@ -94,6 +100,7 @@ namespace ToolFramework{ 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. + Services* m_services; DAQUtilities* m_util; zmq::context_t* m_context; zmq::socket_t* m_pub;