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;