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 : }
|