Line data Source code
1 : /*
2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
3 : */
4 :
5 : #include <boost/assign/list_of.hpp>
6 : #include "base/regex.h"
7 : #include "rapidjson/document.h"
8 : #include "rapidjson/stringbuffer.h"
9 : #include "rapidjson/writer.h"
10 : #include "query.h"
11 : #include "stats_query.h"
12 :
13 : using boost::assign::map_list_of;
14 : using contrail::regex;
15 : using contrail::regex_match;
16 : using contrail::regex_search;
17 :
18 1811 : bool PostProcessingQuery::sort_field_comparator(
19 : const QEOpServerProxy::ResultRowT& lhs,
20 : const QEOpServerProxy::ResultRowT& rhs) {
21 1811 : std::map<std::string, std::string>::const_iterator lhs_it, rhs_it;
22 1811 : for (std::vector<sort_field_t>::iterator sort_it = sort_fields.begin();
23 2864 : sort_it != sort_fields.end(); sort_it++) {
24 1811 : lhs_it = lhs.first.find((*sort_it).name);
25 1811 : QE_ASSERT(lhs_it != lhs.first.end());
26 1811 : rhs_it = rhs.first.find((*sort_it).name);
27 1811 : QE_ASSERT(rhs_it != rhs.first.end());
28 7244 : if ((*sort_it).type == std::string("int") ||
29 9055 : (*sort_it).type == std::string("long") ||
30 3622 : (*sort_it).type == std::string("ipv4")) {
31 0 : uint64_t lhs_val = 0, rhs_val = 0;
32 0 : stringToInteger(lhs_it->second, lhs_val);
33 0 : stringToInteger(rhs_it->second, rhs_val);
34 0 : if (lhs_val < rhs_val) return true;
35 0 : if (lhs_val > rhs_val) return false;
36 : } else {
37 1811 : if (lhs_it->second < rhs_it->second) return true;
38 1193 : if (lhs_it->second > rhs_it->second) return false;
39 : }
40 : }
41 :
42 1053 : return false;
43 : }
44 :
45 32 : bool PostProcessingQuery::merge_processing(
46 : const QEOpServerProxy::BufferT& input,
47 : QEOpServerProxy::BufferT& output)
48 : {
49 32 : if (status_details != 0)
50 : {
51 0 : QE_TRACE(DEBUG,
52 : "No need to process query, as there were errors previously");
53 0 : return false;
54 : }
55 :
56 : // Check if the result has to be sorted
57 32 : if (sorted) {
58 24 : QEOpServerProxy::BufferT *merged_result = &output;
59 24 : const QEOpServerProxy::BufferT *raw_result1 = &(input);
60 :
61 24 : if (result_.get() == NULL) {
62 24 : size_t merged_result_size = merged_result->size();
63 24 : merged_result->reserve(merged_result_size + raw_result1->size());
64 24 : copy(raw_result1->begin(), raw_result1->end(),
65 : std::back_inserter(*merged_result));
66 24 : if (merged_result_size) {
67 0 : if (sorting_type == ASCENDING) {
68 0 : std::inplace_merge(merged_result->begin(),
69 0 : merged_result->begin() + merged_result_size,
70 : merged_result->end(),
71 : boost::bind(&PostProcessingQuery::sort_field_comparator,
72 : this, _1, _2));
73 : } else {
74 0 : std::inplace_merge(merged_result->rbegin(),
75 0 : merged_result->rbegin() + raw_result1->size(),
76 0 : merged_result->rend(),
77 : boost::bind(&PostProcessingQuery::sort_field_comparator,
78 : this, _1, _2));
79 : }
80 : }
81 : } else {
82 0 : QEOpServerProxy::BufferT *raw_result2 = result_.get();
83 0 : size_t size1 = raw_result1->size();
84 0 : size_t size2 = raw_result2->size();
85 0 : QE_TRACE(DEBUG, "Merging results from vectors of size:" <<
86 : size1 << " and " << size2);
87 0 : merged_result->reserve(raw_result1->size() + raw_result2->size());
88 0 : if (sorting_type == ASCENDING) {
89 0 : std::merge(raw_result1->begin(), raw_result1->end(),
90 : raw_result2->begin(), raw_result2->end(),
91 : std::back_inserter(*merged_result),
92 : boost::bind(&PostProcessingQuery::sort_field_comparator,
93 : this, _1, _2));
94 : } else {
95 0 : std::merge(raw_result1->rbegin(), raw_result1->rend(),
96 0 : raw_result2->rbegin(), raw_result2->rend(),
97 : std::back_inserter(*merged_result),
98 : boost::bind(&PostProcessingQuery::sort_field_comparator,
99 : this, _1, _2));
100 : }
101 : }
102 : } else {
103 8 : QE_TRACE(DEBUG, "Merge_Processing: Adding inputs to output");
104 8 : QEOpServerProxy::BufferT *merged_result = &output;
105 8 : const QEOpServerProxy::BufferT *raw_result1 = &(input);
106 :
107 8 : if (result_.get() == NULL)
108 : {
109 8 : merged_result->reserve(raw_result1->size());
110 8 : copy(raw_result1->begin(), raw_result1->end(),
111 : std::back_inserter(*merged_result));
112 : } else {
113 :
114 0 : QEOpServerProxy::BufferT *raw_result2 = result_.get();
115 0 : size_t size1 = raw_result1->size();
116 0 : size_t size2 = raw_result2->size();
117 0 : QE_TRACE(DEBUG, "Merging results from vectors of size:" <<
118 : size1 << " and " << size2);
119 0 : merged_result->reserve(raw_result1->size() + raw_result2->size());
120 0 : copy(raw_result1->begin(), raw_result1->end(),
121 : std::back_inserter(*merged_result));
122 0 : copy(raw_result2->begin(), raw_result2->end(),
123 : std::back_inserter(*merged_result));
124 : }
125 8 : QE_TRACE(DEBUG, "Merge_Processing: Done adding inputs to output");
126 : }
127 :
128 : // Have the result ready and processing is done
129 32 : status_details = 0;
130 32 : return true;
131 : }
132 :
133 4 : bool PostProcessingQuery::final_merge_processing(
134 : const std::vector<boost::shared_ptr<QEOpServerProxy::BufferT> >& inputs,
135 : QEOpServerProxy::BufferT& output)
136 : {
137 4 : bool merge_done = false;
138 :
139 4 : if (status_details != 0)
140 : {
141 0 : QE_TRACE(DEBUG,
142 : "No need to process query, as there were errors previously");
143 0 : return false;
144 : }
145 :
146 4 : if (!merge_done) {
147 4 : QEOpServerProxy::BufferT *merged_result = &output;
148 4 : size_t final_vector_size = 0;
149 : // merge the results from parallel queries
150 36 : for (size_t i = 0; i < inputs.size(); i++) {
151 32 : final_vector_size += inputs[i]->size();
152 : }
153 4 : merged_result->reserve(final_vector_size);
154 4 : QE_TRACE(DEBUG, "Merging results between " << inputs.size()
155 : << " vectors with final vector size:" << final_vector_size);
156 36 : for (size_t i = 0; i < inputs.size(); i++) {
157 32 : QEOpServerProxy::BufferT *raw_result = inputs[i].get();
158 32 : copy(raw_result->begin(), raw_result->end(),
159 : std::back_inserter(*merged_result));
160 : }
161 : }
162 :
163 4 : if (sorted) {
164 3 : QEOpServerProxy::BufferT *merged_result = &output;
165 3 : if (sorting_type == ASCENDING) {
166 2 : std::sort(merged_result->begin(), merged_result->end(),
167 : boost::bind(&PostProcessingQuery::sort_field_comparator,
168 : this, _1, _2));
169 : } else {
170 1 : std::sort(merged_result->rbegin(), merged_result->rend(),
171 : boost::bind(&PostProcessingQuery::sort_field_comparator,
172 : this, _1, _2));
173 : }
174 : }
175 :
176 4 : if (limit) {
177 2 : QEOpServerProxy::BufferT *merged_result = &output;
178 2 : QE_TRACE(DEBUG, "Apply Limit [" << limit << "]");
179 2 : if (merged_result->size() > (size_t)limit) {
180 0 : merged_result->resize(limit);
181 : }
182 : }
183 :
184 : // Have the result ready and processing is done
185 4 : status_details = 0;
186 4 : return true;
187 : }
188 :
189 865 : query_status_t PostProcessingQuery::process_query() {
190 865 : if (status_details != 0)
191 : {
192 0 : QE_TRACE(DEBUG,
193 : "No need to process query, as there were errors previously");
194 0 : return QUERY_FAILURE;
195 : }
196 :
197 865 : AnalyticsQuery *mquery = (AnalyticsQuery *)main_query;
198 865 : result_ = std::move(mquery->selectquery_->result_);
199 865 : mresult_ = std::move(mquery->selectquery_->mresult_);
200 865 : QEOpServerProxy::BufferT *raw_result = result_.get();
201 :
202 : /* filter are ANDs over OR
203 : * [ [ e1 AND e2 ] OR [ e3 ] ]
204 : */
205 : /* below is filter processing for stats table queries
206 : */
207 865 : if (filter_list.size()) {
208 244 : size_t num_filtered=0;
209 244 : MapBufT::iterator kt = mresult_->end();
210 244 : for (MapBufT::iterator it = mresult_->begin();
211 376 : it!= mresult_->end(); it++) {
212 :
213 132 : if (kt!=mresult_->end()) {
214 83 : mresult_->erase(kt);
215 83 : kt = mresult_->end();
216 : }
217 132 : std::map<std::string, QEOpServerProxy::SubVal>& attrs = it->second.first;
218 132 : bool delete_row = true;
219 132 : std::string unknown_attr;
220 230 : for (size_t j = 0; j < filter_list.size(); j++) {
221 132 : std::vector<filter_match_t>& filter_and = filter_list[j];
222 132 : bool and_check = true;
223 :
224 160 : for (size_t k = 0; k < filter_and.size(); k++) {
225 : std::map<std::string, QEOpServerProxy::SubVal>::const_iterator iter =
226 132 : attrs.find(filter_and[k].name);
227 132 : if (iter == attrs.end()) {
228 6 : unknown_attr = filter_and[k].name;
229 104 : break;
230 : } else {
231 126 : unknown_attr.clear();
232 : }
233 126 : std::ostringstream vstream;
234 126 : vstream << iter->second;
235 :
236 126 : switch(filter_and[k].op) {
237 0 : case EQUAL:
238 0 : if (filter_and[k].value != vstream.str())
239 : {
240 0 : and_check = false;
241 : }
242 0 : break;
243 0 : case NOT_EQUAL:
244 0 : if (filter_and[k].value == vstream.str())
245 : {
246 0 : and_check = false;
247 : }
248 0 : break;
249 8 : case LEQ:
250 8 : if (iter->second.which() ==
251 : QEOpServerProxy::UINT64) {
252 : uint64_t filter_val;
253 4 : stringToInteger(filter_and[k].value,
254 : filter_val);
255 : uint64_t col_val =
256 4 : boost::get<uint64_t>(iter->second);
257 4 : if (col_val > filter_val) {
258 2 : and_check = false;
259 : }
260 4 : } else if (iter->second.which() ==
261 : QEOpServerProxy::DOUBLE) {
262 : double filter_val;
263 4 : stringToInteger(filter_and[k].value,
264 : filter_val);
265 : double col_val =
266 4 : boost::get<double>(iter->second);
267 4 : if (col_val > filter_val) {
268 3 : and_check = false;
269 : }
270 : }
271 8 : break;
272 8 : case GEQ:
273 8 : if (iter->second.which() ==
274 : QEOpServerProxy::UINT64) {
275 : uint64_t filter_val;
276 4 : stringToInteger(filter_and[k].value,
277 : filter_val);
278 : uint64_t col_val =
279 4 : boost::get<uint64_t>(iter->second);
280 4 : if (col_val < filter_val) {
281 2 : and_check = false;
282 : }
283 4 : } else if (iter->second.which() ==
284 : QEOpServerProxy::DOUBLE) {
285 : double filter_val;
286 4 : stringToInteger(filter_and[k].value,
287 : filter_val);
288 : double col_val =
289 4 : boost::get<double>(iter->second);
290 4 : if (col_val < filter_val) {
291 2 : and_check = false;
292 : }
293 : }
294 8 : break;
295 110 : case REGEX_MATCH:
296 110 : if (!regex_match(vstream.str(),
297 110 : filter_and[k].match_e)) {
298 89 : and_check = false;
299 : }
300 110 : break;
301 0 : default:
302 : // upsupported filter operation
303 0 : QE_TRACE(DEBUG, "Unsupported filter operation " << filter_and[k].op);
304 0 : break;
305 : }
306 126 : if (and_check == false)
307 98 : break;
308 126 : }
309 132 : if (and_check == true) {
310 34 : delete_row = false;
311 34 : break;
312 : }
313 : }
314 132 : if (unknown_attr.size()) {
315 6 : QE_TRACE(DEBUG, "Unknown filter attr in row " << unknown_attr);
316 : }
317 132 : if (delete_row) {
318 98 : num_filtered++;
319 98 : kt = it;
320 : }
321 132 : }
322 244 : if (kt!=mresult_->end()) {
323 15 : mresult_->erase(kt);
324 15 : kt = mresult_->end();
325 : }
326 244 : QE_TRACE(DEBUG, "# of entries filtered is " << num_filtered);
327 :
328 : }
329 :
330 : /* below is filter processing for non stats table queries
331 : */
332 865 : if (filter_list.size() != 0) {
333 244 : QEOpServerProxy::BufferT filtered_table;
334 : // do filter operation
335 244 : QE_TRACE(DEBUG, "Doing filter operation");
336 1332 : for (size_t i = 0; i < raw_result->size(); i++) {
337 1088 : QEOpServerProxy::ResultRowT row = (*raw_result)[i];
338 1088 : bool delete_row = true;
339 :
340 1468 : for (size_t j = 0; j < filter_list.size(); j++) {
341 1127 : std::vector<filter_match_t>& filter_and = filter_list[j];
342 1127 : bool and_check = true;
343 :
344 2026 : for (size_t k = 0; k < filter_and.size(); k++) {
345 1279 : std::map<std::string, std::string>::iterator iter;
346 1279 : iter = row.first.find(filter_and[k].name);
347 1279 : if (iter == row.first.end())
348 : {
349 553 : if (!(filter_and[k].ignore_col_absence)) {
350 0 : and_check = false;
351 380 : break;
352 : }
353 553 : continue;
354 : }
355 :
356 726 : switch(filter_and[k].op)
357 : {
358 180 : case EQUAL:
359 180 : if (filter_and[k].value != iter->second)
360 : {
361 102 : and_check = false;
362 : }
363 180 : break;
364 :
365 221 : case NOT_EQUAL:
366 221 : if (filter_and[k].value == iter->second)
367 : {
368 49 : and_check = false;
369 : }
370 221 : break;
371 :
372 0 : case LEQ:
373 : {
374 : int filter_value =
375 0 : atoi(filter_and[k].value.c_str());
376 0 : int column_value= atoi(iter->second.c_str());
377 0 : if (column_value > filter_value)
378 : {
379 0 : and_check = false;
380 : }
381 0 : break;
382 : }
383 :
384 0 : case GEQ:
385 : {
386 : int filter_value =
387 0 : atoi(filter_and[k].value.c_str());
388 0 : int column_value= atoi(iter->second.c_str());
389 0 : if (column_value < filter_value)
390 : {
391 0 : and_check = false;
392 : }
393 0 : break;
394 : }
395 :
396 325 : case REGEX_MATCH:
397 : {
398 325 : if (!regex_match(iter->second,
399 325 : filter_and[k].match_e))
400 : {
401 229 : and_check = false;
402 : }
403 325 : break;
404 : }
405 :
406 0 : default:
407 : // upsupported filter operation
408 0 : QE_LOG(ERROR, "Unsupported filter operation: " <<
409 : filter_and[k].op);
410 0 : return QUERY_FAILURE;
411 : }
412 726 : if (and_check == false)
413 380 : break;
414 : }
415 :
416 1127 : if (and_check == true) {
417 747 : QE_TRACE(DEBUG, "filter out entry #:" << i);
418 747 : delete_row = false;
419 747 : break;
420 : }
421 : }
422 1088 : if (!delete_row) {
423 747 : filtered_table.push_back(row);
424 : }
425 1088 : }
426 244 : *raw_result = filtered_table;
427 244 : }
428 :
429 : // Check if the result has to be sorted
430 865 : if (sorted) {
431 34 : if (sorting_type == ASCENDING) {
432 26 : std::sort(raw_result->begin(), raw_result->end(),
433 : boost::bind(&PostProcessingQuery::sort_field_comparator,
434 : this, _1, _2));
435 : } else {
436 8 : std::sort(raw_result->rbegin(), raw_result->rend(),
437 : boost::bind(&PostProcessingQuery::sort_field_comparator,
438 : this, _1, _2));
439 : }
440 : }
441 :
442 : // If the flow series query is parallelized, we should apply the limit
443 : // only after the result from all the tasks are merged
444 : // (@ final_merge_processing).
445 880 : if ((mquery->table() != g_viz_constants.FLOW_SERIES_TABLE ||
446 880 : (mquery->table() == g_viz_constants.FLOW_SERIES_TABLE &&
447 880 : !mquery->is_query_parallelized())) && limit) {
448 27 : QE_TRACE(DEBUG, "Apply Limit [" << limit << "]");
449 27 : if (raw_result->size() > (size_t)limit) {
450 2 : raw_result->resize(limit);
451 : }
452 27 : if (mresult_->size() > (size_t)limit) {
453 4 : MapBufT::iterator it = mresult_->begin();
454 4 : std::advance(it, limit);
455 4 : mresult_->erase(it, mresult_->end());
456 : }
457 : }
458 :
459 865 : if (IS_TRACE_ENABLED(POSTPROCESS_RESULT_TRACE))
460 : {
461 0 : std::vector<QEOpServerProxy::ResultRowT>::iterator res_it;
462 0 : QE_TRACE(DEBUG, "== Post Processing Result ==");
463 0 : for (res_it = raw_result->begin(); res_it != raw_result->end();
464 0 : ++res_it) {
465 0 : std::vector<final_result_col> row_entry;
466 0 : std::map<std::string, std::string>::iterator map_it;
467 0 : for (map_it = (*res_it).first.begin();
468 0 : map_it != (*res_it).first.end(); ++map_it) {
469 0 : final_result_col col;
470 0 : col.set_col(map_it->first); col.set_value(map_it->second);
471 0 : row_entry.push_back(col);
472 : //QE_TRACE(DEBUG, map_it->first << " : " << map_it->second);
473 0 : }
474 0 : FINAL_RESULT_ROW_TRACE(QeTraceBuf, mquery->query_id, row_entry);
475 0 : }
476 : }
477 :
478 : #if 0
479 : //if ((limit) && (!sorted))
480 : for (int i = 0 ; i < 200000; i++)
481 : raw_result->push_back(map_list_of(
482 : "destvn","abc-cor\"poration:front-end-network:001")(
483 : "sourceip","168430090")("destip","3232238090")(
484 : "sourcevn","abc-corporation:front-end-network:002")(
485 : "protocol","80")("dport","62000")("sport","1000")(
486 : "sum(packets)","4294967196")
487 : );
488 : #endif
489 :
490 : // Have the result ready and processing is done
491 865 : status_details = 0;
492 865 : parent_query->subquery_processed(this);
493 865 : return QUERY_SUCCESS;
494 : }
|