LCOV - code coverage report
Current view: top level - root/contrail/src/contrail-analytics/contrail-collector - generator.cc (source / functions) Hit Total Coverage
Test: OpenSDN C/C++ coverage (all TARGET_SET jobs) Lines: 13 229 5.7 %
Date: 2026-08-03 02:19:58 Functions: 1 35 2.9 %
Legend: Lines: hit not hit

          Line data    Source code
       1             : /*
       2             :  * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
       3             :  */
       4             : 
       5             : #include <utility>
       6             : #include <vector>
       7             : #include <map>
       8             : #include <boost/bind/bind.hpp>
       9             : #include "base/timer.h"
      10             : #include <boost/date_time/posix_time/posix_time.hpp>
      11             : #include <boost/assign/list_of.hpp>
      12             : 
      13             : #include <sandesh/sandesh_trace.h>
      14             : #include <sandesh/sandesh_session.h>
      15             : #include <sandesh/sandesh_ctrl_types.h>
      16             : #include <sandesh/common/vns_types.h>
      17             : #include <sandesh/common/vns_constants.h>
      18             : #include <sandesh/sandesh_uve_types.h>
      19             : #include <sandesh/sandesh_message_builder.h>
      20             : 
      21             : #include "OpServerProxy.h"
      22             : #include "db_handler.h"
      23             : #include "collector.h"
      24             : #include "generator.h"
      25             : #include "viz_collector.h"
      26             : #include "viz_sandesh.h"
      27             : #include <analytics/viz_types.h>
      28             : #include "vizd_table_desc.h"
      29             : 
      30             : extern SandeshTraceBufferPtr UVETraceBuf;
      31             : 
      32             : using std::string;
      33             : using std::pair;
      34             : using std::vector;
      35             : using std::map;
      36             : using namespace boost::placeholders;
      37             : 
      38             : #define GENERATOR_LOG(_Level, _Msg)                                            \
      39             :     do {                                                                       \
      40             :         if (LoggingDisabled()) break;                                          \
      41             :         log4cplus::Logger _Xlogger = log4cplus::Logger::getRoot();             \
      42             :         if (_Xlogger.isEnabledFor(log4cplus::_Level##_LOG_LEVEL)) {            \
      43             :             log4cplus::tostringstream _Xbuf;                                   \
      44             :             _Xbuf << ToString() << ": " << __func__ << ": " << _Msg;           \
      45             :             _Xlogger.forcedLog(log4cplus::_Level##_LOG_LEVEL,                  \
      46             :                                _Xbuf.str());                                   \
      47             :         }                                                                      \
      48             :     } while (false)
      49             : 
      50           0 : void Generator::UpdateStatistics(const VizMsg *vmsg) {
      51           0 :     std::scoped_lock lock(smutex_);
      52           0 :     statistics_.Update(vmsg);
      53           0 : }
      54             : 
      55           0 : void Generator::GetStatistics(vector<SandeshStats> *ssv) const {
      56           0 :     std::scoped_lock lock(smutex_);
      57           0 :     statistics_.Get(ssv);
      58           0 : }
      59             : 
      60           0 : void Generator::GetStatistics(vector<SandeshLogLevelStats> *lsv) const {
      61           0 :     std::scoped_lock lock(smutex_);
      62           0 :     statistics_.Get(lsv);
      63           0 : }
      64             : 
      65           0 : void Generator::SendSandeshMessageStatistics() {
      66           0 :     vector<SandeshMessageInfo> smv;
      67             :     {
      68           0 :         std::scoped_lock lock(smutex_);
      69           0 :         statistics_.Get(&smv);
      70           0 :     }
      71           0 :     SandeshMessageStat * snh = SANDESH_MESSAGE_STAT_CREATE();
      72           0 :     snh->set_name(ToString());
      73           0 :     snh->set_msg_info(smv);
      74           0 :     SANDESH_MESSAGE_STAT_SEND_SANDESH(snh);
      75           0 : }
      76             :     
      77           0 : bool Generator::ReceiveSandeshMsg(const VizMsg *vmsg, bool rsc) {
      78           0 :     UpdateStatistics(vmsg);
      79           0 :     return ProcessRules(vmsg, rsc);
      80             : }
      81             : 
      82             : // SandeshGenerator
      83           0 : SandeshGenerator::SandeshGenerator(Collector * const collector, VizSession *session,
      84             :         SandeshStateMachine *state_machine, const string &source,
      85             :         const string &module, const string &instance_id,
      86             :         const string &node_type,
      87           0 :         DbHandlerPtr global_db_handler) :
      88             :         Generator(),
      89           0 :         collector_(collector),
      90           0 :         state_machine_(state_machine),
      91           0 :         viz_session_(session),
      92           0 :         instance_id_(instance_id),
      93           0 :         node_type_(node_type),
      94           0 :         source_(source),
      95           0 :         module_(module),
      96           0 :         name_(source + ":" + node_type_ + ":" + module + ":" + instance_id_),
      97           0 :         instance_(session->GetSessionInstance()),
      98           0 :         process_rules_cb_(
      99             :             boost::bind(&SandeshGenerator::ProcessRulesCb, this, _1)),
     100           0 :         sm_defer_timer_(NULL),
     101           0 :         sm_defer_timer_expiry_time_usec_(0),
     102           0 :         sm_defer_time_msec_(0) {
     103             :         //Use collector db_handler
     104           0 :         db_handler_ = global_db_handler;
     105           0 :         disconnected_ = false;
     106           0 :         gen_attr_.set_connects(1);
     107           0 :         gen_attr_.set_connect_time(UTCTimestampUsec());
     108             :         // Update state machine
     109           0 :         state_machine_->SetGeneratorKey(name_);
     110           0 :         CreateStateMachineDeferTimer();
     111           0 : }
     112             : 
     113           0 : SandeshGenerator::~SandeshGenerator() {
     114           0 :     DeleteStateMachineDeferTimer();
     115           0 : }
     116             : 
     117           0 : void SandeshGenerator::set_session(VizSession *session) {
     118           0 :     viz_session_ = session;
     119           0 :     instance_ = session->GetSessionInstance();
     120           0 :     session->set_generator(this);
     121           0 : }
     122             : 
     123           0 : void SandeshGenerator::TimerErrorHandler(string name, string error) {
     124           0 :     GENERATOR_LOG(ERROR, name + " error: " + error);
     125           0 : }
     126             : 
     127           0 : void SandeshGenerator::ReceiveSandeshCtrlMsg(uint32_t connects) {
     128             :     // This is a control message during SandeshGenerator-Collector negotiation
     129           0 :     ModuleServerState ginfo;
     130           0 :     GetGeneratorInfo(ginfo);
     131           0 :     SandeshModuleServerTrace::Send(ginfo);
     132             :     // Setup state machine watermarks
     133           0 :     std::vector<Sandesh::QueueWaterMarkInfo> wm_info;
     134           0 :     collector_->GetSmQueueWaterMarkInfo(wm_info);
     135           0 :     for (size_t i = 0; i < wm_info.size(); i++) {
     136           0 :         state_machine_->SetQueueWaterMarkInfo(wm_info[i]);
     137             :     }
     138           0 : }
     139             : 
     140           0 : void SandeshGenerator::DisconnectSession(VizSession *vsession) {
     141           0 :     std::scoped_lock lock(mutex_);
     142           0 :     GENERATOR_LOG(INFO, "Session:" << vsession->ToString());
     143           0 :     if (vsession == viz_session_) {
     144           0 :         disconnected_ = true;
     145             :         // This SandeshGenerator's session is now gone.
     146             :         // Delete all its UVEs
     147           0 :         uint32_t tmp = gen_attr_.get_resets();
     148           0 :         gen_attr_.set_resets(tmp+1);
     149           0 :         gen_attr_.set_reset_time(UTCTimestampUsec());
     150           0 :         state_machine_->ResetQueueWaterMarkInfo();
     151           0 :         StopStateMachineDeferTimer();
     152           0 :         DeleteStateMachineDeferTimer();
     153           0 :         sm_defer_timer_expiry_time_usec_ = 0;
     154           0 :         sm_defer_time_msec_ = 0;
     155           0 :         viz_session_ = NULL;
     156           0 :         state_machine_ = NULL;
     157           0 :         vsession->set_generator(NULL);
     158           0 :         collector_->GetOSP()->DeleteUVEs(source_, module_, 
     159           0 :                                          node_type_, instance_id_);
     160           0 :         ModuleServerState ginfo;
     161           0 :         GetGeneratorInfo(ginfo);
     162           0 :         SandeshModuleServerTrace::Send(ginfo);
     163           0 :     } else {
     164           0 :         GENERATOR_LOG(ERROR, "Disconnect for session:" << vsession->ToString() <<
     165             :                 ", generator session:" << viz_session_->ToString());
     166             :     }
     167           0 : }
     168             : 
     169           0 : bool SandeshGenerator::StateMachineDeferTimerExpired() {
     170           0 :     std::scoped_lock lock(mutex_);
     171           0 :     sm_defer_timer_expiry_time_usec_ = UTCTimestampUsec();
     172           0 :     if (state_machine_) {
     173           0 :         state_machine_->SetDeferDequeue(false);
     174             :     }
     175           0 :     return false;
     176           0 : }
     177             : 
     178           0 : void SandeshGenerator::CreateStateMachineDeferTimer() {
     179             :     // Run in the context of sandesh state machine task
     180           0 :     assert(sm_defer_timer_ == NULL);
     181           0 :     sm_defer_timer_ = TimerManager::CreateTimer(
     182           0 :         *collector_->event_manager()->io_service(),
     183           0 :         "SandeshGenerator SM Defer Timer: " + name_,
     184           0 :         state_machine_->connection()->GetTaskId(), instance_);
     185           0 : }
     186             : 
     187           0 : void SandeshGenerator::StartStateMachineDeferTimer(int time_msec) {
     188           0 :     sm_defer_timer_->Start(time_msec,
     189             :             boost::bind(
     190             :                 &SandeshGenerator::StateMachineDeferTimerExpired, this),
     191             :             boost::bind(&SandeshGenerator::TimerErrorHandler, this, _1, _2));
     192           0 : }
     193             : 
     194           0 : void SandeshGenerator::StopStateMachineDeferTimer() {
     195           0 :     assert(sm_defer_timer_->Cancel());
     196           0 : }
     197             : 
     198           0 : void SandeshGenerator::DeleteStateMachineDeferTimer() {
     199           0 :     TimerManager::DeleteTimer(sm_defer_timer_);
     200           0 :     sm_defer_timer_ = NULL;
     201           0 : }
     202             : 
     203           0 : bool SandeshGenerator::IsStateMachineDeferTimerRunningUnlocked() const {
     204           0 :     if (sm_defer_timer_) {
     205           0 :         return sm_defer_timer_->running();
     206             :     }
     207           0 :     return false;
     208             : }
     209             : 
     210           0 : bool SandeshGenerator::IsStateMachineDeferTimerRunning() const {
     211           0 :     std::scoped_lock lock(mutex_);
     212           0 :     return IsStateMachineDeferTimerRunningUnlocked();
     213           0 : }
     214             : 
     215           0 : int SandeshGenerator::GetStateMachineDeferTimeMSec() const {
     216           0 :     std::scoped_lock lock(mutex_);
     217           0 :     return sm_defer_time_msec_;
     218           0 : }
     219             : 
     220           7 : int GetDeferTimeMSec(uint64_t event_time_usec,
     221             :     uint64_t last_expiry_time_usec, uint64_t last_defer_time_usec) {
     222             :     // If this is the first time, then defer the state machine with
     223             :     // initial defer time
     224           7 :     if (last_defer_time_usec == 0 || last_expiry_time_usec == 0) {
     225           2 :         return SandeshGenerator::kInitialSmDeferTimeMSec;
     226             :     }
     227           5 :     assert(event_time_usec >= last_expiry_time_usec);
     228           5 :     uint64_t time_since_expiry_usec(event_time_usec - last_expiry_time_usec);
     229             :     // We will double the defer time if we get a back pressure
     230             :     // event within 2 * last defer time. If the back pressure
     231             :     // event is between 2 * last defer time and 4 * last defer
     232             :     // time, then the defer time will be same as the current
     233             :     // defer time. If the back pressure event is after  4 * last
     234             :     // defer time, then we will reset the defer time to the
     235             :     // initial defer time
     236           5 :     if (time_since_expiry_usec <= 2 * last_defer_time_usec) {
     237           3 :         uint64_t ndefer_time_msec((2 * last_defer_time_usec)/1000);
     238           6 :         return std::min(ndefer_time_msec,
     239           3 :             static_cast<uint64_t>(SandeshGenerator::kMaxSmDeferTimeMSec));
     240           2 :     } else if ((2 * last_defer_time_usec <= time_since_expiry_usec) &&
     241           2 :         (time_since_expiry_usec <= 4 * last_defer_time_usec)) {
     242           1 :         return last_defer_time_usec/1000;
     243             :     } else {
     244           1 :         return SandeshGenerator::kInitialSmDeferTimeMSec;
     245             :     }
     246             : }
     247             : 
     248           0 : void SandeshGenerator::ProcessRulesCb(GenDb::DbOpResult::type dresult) {
     249           0 :     std::scoped_lock lock(mutex_);
     250           0 :     if (dresult == GenDb::DbOpResult::BACK_PRESSURE) {
     251           0 :         if (state_machine_) {
     252             :             // If state mchine defer timer is running just return to
     253             :             // avoid increasing the defer time more than once every
     254             :             // timer expiry
     255           0 :             if (IsStateMachineDeferTimerRunningUnlocked()) {
     256           0 :                 return;
     257             :             }
     258           0 :             state_machine_->SetDeferDequeue(true);
     259           0 :             uint64_t now_usec(UTCTimestampUsec());
     260           0 :             int defer_time_msec(GetDeferTimeMSec(now_usec,
     261           0 :                 sm_defer_timer_expiry_time_usec_, sm_defer_time_msec_ * 1000));
     262           0 :             sm_defer_time_msec_ = defer_time_msec;
     263           0 :             StartStateMachineDeferTimer(sm_defer_time_msec_);
     264             :         }
     265             :     }
     266           0 : }
     267             : 
     268           0 : bool SandeshGenerator::ProcessRules(const VizMsg *vmsg, bool rsc) {
     269           0 :     return collector_->ProcessSandeshMsgCb()(vmsg, rsc, GetDbHandler(),
     270           0 :         process_rules_cb_);
     271             : }
     272             : 
     273           0 : bool SandeshGenerator::GetSandeshStateMachineQueueCount(
     274             :     uint64_t &queue_count) const {
     275           0 :     if (!state_machine_) {
     276             :         // Return 0 so that last stale value is not displayed
     277           0 :         queue_count = 0;
     278           0 :         return true;
     279             :     }
     280           0 :     return state_machine_->GetQueueCount(queue_count);
     281             : }
     282             : 
     283           0 : bool SandeshGenerator::GetSandeshStateMachineDropLevel(
     284             :     std::string &drop_level) const {
     285           0 :     if (!state_machine_) {
     286           0 :         return false;
     287             :     }
     288           0 :     return state_machine_->GetMessageDropLevel(drop_level);
     289             : }
     290             : 
     291           0 : bool SandeshGenerator::GetSandeshStateMachineStats(
     292             :                     SandeshStateMachineStats &sm_stats,
     293             :                     SandeshGeneratorBasicStats &sm_msg_stats) const {
     294           0 :     if (!state_machine_) {
     295           0 :         return false;
     296             :     }
     297           0 :     return state_machine_->GetStatistics(sm_stats, sm_msg_stats);
     298             : }
     299             : 
     300           0 : void SandeshGenerator::GetGeneratorInfo(ModuleServerState &genlist) const {
     301           0 :     vector<GeneratorInfo> giv;
     302           0 :     GeneratorInfo gi;
     303           0 :     gi.set_hostname(Sandesh::source());
     304           0 :     gi.set_gen_attr(gen_attr_);
     305           0 :     giv.push_back(gi);
     306           0 :     genlist.set_generator_info(giv);
     307           0 :     genlist.set_name(source() + ":" + node_type_ + ":" + module() + ":" +
     308           0 :         instance_id_);
     309           0 : }
     310             : 
     311           0 : const std::string SandeshGenerator::State() const {
     312           0 :     if (state_machine_) {
     313           0 :         return state_machine_->StateName();
     314             :     }
     315           0 :     return "Disconnected";
     316             : }
     317             : 
     318           0 : void SandeshGenerator::ConnectSession(VizSession *session,
     319             :     SandeshStateMachine *state_machine) {
     320           0 :     std::scoped_lock lock(mutex_);
     321           0 :     set_session(session);
     322           0 :     set_state_machine(state_machine);
     323           0 :     disconnected_ = false;
     324           0 :     uint32_t tmp = gen_attr_.get_connects();
     325           0 :     gen_attr_.set_connects(tmp+1);
     326           0 :     gen_attr_.set_connect_time(UTCTimestampUsec());
     327           0 :     CreateStateMachineDeferTimer();
     328           0 : }
     329             : 
     330           0 : void SandeshGenerator::SetDbQueueWaterMarkInfo(
     331             :     Sandesh::QueueWaterMarkInfo &wm) {
     332           0 :     if (!GetDbHandler()) {
     333           0 :         return;
     334             :     }
     335           0 :     bool high(boost::get<2>(wm));
     336           0 :     bool defer_undefer(boost::get<3>(wm));
     337           0 :     boost::function<void (void)> cb;
     338           0 :     if (high && defer_undefer) {
     339           0 :         cb = boost::bind(&SandeshStateMachine::SetDeferDequeue,
     340           0 :                 state_machine_, true);
     341           0 :     } else if (!high && defer_undefer) {
     342             :         cb = boost::bind(&SandeshStateMachine::SetDeferDequeue,
     343           0 :                 state_machine_, false);
     344             :     }
     345           0 :     GetDbHandler()->SetDbQueueWaterMarkInfo(wm, cb);
     346           0 : }
     347             : 
     348           0 : void SandeshGenerator::ResetDbQueueWaterMarkInfo() {
     349           0 :     if (!GetDbHandler()) {
     350           0 :         return;
     351             :     }
     352           0 :     GetDbHandler()->ResetDbQueueWaterMarkInfo();
     353             : }
     354             : 
     355           0 : void SandeshGenerator::SetSmQueueWaterMarkInfo(
     356             :     Sandesh::QueueWaterMarkInfo &wm) {
     357           0 :     if (state_machine_) {
     358           0 :         state_machine_->SetQueueWaterMarkInfo(wm);
     359             :     }
     360           0 : }
     361             : 
     362           0 : void SandeshGenerator::ResetSmQueueWaterMarkInfo() {
     363           0 :     if (state_machine_) {
     364           0 :         state_machine_->ResetQueueWaterMarkInfo();
     365             :     }
     366           0 : }
     367             : 
     368             : // SyslogGenerator
     369           0 : SyslogGenerator::SyslogGenerator(SyslogListeners *const listeners,
     370           0 :         const string &source, const string &module) :
     371             :           Generator(),
     372           0 :           syslog_(listeners),
     373           0 :           source_(source),
     374           0 :           module_(module),
     375           0 :           name_(source + ":" + module),
     376           0 :           db_handler_(listeners->GetDbHandler()) {
     377           0 : }
     378             : 
     379           0 : bool SyslogGenerator::ProcessRules(const VizMsg *vmsg, bool rsc) {
     380           0 :     return syslog_->ProcessSandeshMsgCb()(vmsg, rsc, GetDbHandler(),
     381           0 :         GenDb::GenDbIf::DbAddColumnCb());
     382             : }

Generated by: LCOV version 1.14