LCOV - code coverage report
Current view: top level - root/contrail/src/contrail-analytics/contrail-query-engine - stats_select.cc (source / functions) Hit Total Coverage
Test: OpenSDN C/C++ coverage (all TARGET_SET jobs) Lines: 384 443 86.7 %
Date: 2026-08-03 02:19:58 Functions: 11 11 100.0 %
Legend: Lines: hit not hit

          Line data    Source code
       1             : /*
       2             :  * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
       3             :  */
       4             : #include "viz_constants.h"
       5             : #include "stats_select.h"
       6             : #include "stats_query.h"
       7             : #include "query.h"
       8             : #include <cstdlib>
       9             : #include <boost/assign/list_of.hpp>
      10             : #include <boost/functional/hash.hpp>
      11             : #include "rapidjson/document.h"
      12             : #include "rapidjson/stringbuffer.h"
      13             : #include "rapidjson/writer.h"
      14             : 
      15             : using boost::assign::map_list_of; 
      16             : 
      17             : using std::string;
      18             : using std::map;
      19             : using std::vector;
      20             : using std::set;
      21             : using std::pair;
      22             : using std::make_pair;
      23             : 
      24             : struct Centroid {
      25             : };
      26             : 
      27             : void
      28           6 : StatsSelect::DeleteCentroid (Centroid *t) {
      29           6 :     free(t);
      30           6 : }
      31             : 
      32             : struct TDigest {
      33             : };
      34             : 
      35             : void 
      36           3 : StatsSelect::DeleteTDigest(TDigest *t) {
      37           3 :     TDigest_destroy(t);
      38           3 : }
      39             : 
      40             : bool
      41         196 : StatsSelect::Jsonify(const std::string& table,
      42             :         const std::map<std::string, StatVal>&  uniks,
      43             :         const QEOpServerProxy::AggRowT& aggs, std::string& jstr) {
      44             : 
      45         196 :     contrail_rapidjson::Document dd;
      46         196 :     dd.SetObject();  
      47             : 
      48             :     try {
      49         196 :         for (std::map<std::string, StatVal>::const_iterator it = uniks.begin();
      50         879 :                 it!= uniks.end(); it++) {
      51         683 :             switch (it->second.which()) {
      52         374 :                 case QEOpServerProxy::STRING : {
      53         374 :                         string mapit = boost::get<string>(it->second);
      54         374 :                         contrail_rapidjson::Value val(contrail_rapidjson::kStringType);
      55         374 :                         val.SetString(mapit.c_str(), dd.GetAllocator());
      56         374 :                         contrail_rapidjson::Value vk;
      57         374 :                         dd.AddMember(vk.SetString(it->first.c_str(),
      58             :                                                   dd.GetAllocator()),
      59             :                                      val, dd.GetAllocator());
      60         374 :                     }
      61         374 :                     break;
      62          49 :                 case QEOpServerProxy::UUID : {
      63          49 :                         boost::uuids::uuid mapit = boost::get<boost::uuids::uuid>(it->second);
      64          49 :                         contrail_rapidjson::Value val(contrail_rapidjson::kStringType);
      65          49 :                         std::string ustr = to_string(mapit);
      66          49 :                         val.SetString(ustr.c_str(), dd.GetAllocator());
      67          49 :                         contrail_rapidjson::Value vk;
      68          49 :                         dd.AddMember(vk.SetString(it->first.c_str(),
      69             :                                                   dd.GetAllocator()),
      70             :                                      val, dd.GetAllocator()); 
      71          49 :                     }
      72          49 :                     break;
      73         233 :                 case QEOpServerProxy::UINT64 : {
      74         233 :                         uint64_t mapit = boost::get<uint64_t>(it->second);
      75         233 :                         contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
      76         233 :                         val.SetUint64(mapit);
      77         233 :                         contrail_rapidjson::Value vk;
      78         233 :                         dd.AddMember(vk.SetString(it->first.c_str(),
      79             :                                                   dd.GetAllocator()),
      80             :                                      val, dd.GetAllocator());
      81         233 :                     }
      82         233 :                     break;
      83          27 :                 case QEOpServerProxy::DOUBLE : {
      84          27 :                         double mapit = boost::get<double>(it->second);
      85          27 :                         contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
      86          27 :                         val.SetDouble(mapit);
      87          27 :                         contrail_rapidjson::Value vk;
      88          27 :                         dd.AddMember(vk.SetString(it->first.c_str(),
      89             :                                                   dd.GetAllocator()),
      90             :                                      val, dd.GetAllocator());
      91          27 :                     }
      92          27 :                     break;
      93           0 :                 default:
      94           0 :                     QE_ASSERT(0);
      95             :             }
      96             :         }
      97         196 :         for (QEOpServerProxy::AggRowT::const_iterator it = aggs.begin();
      98         303 :                 it!= aggs.end(); it++) {
      99             :             
     100         107 :             string sname("");
     101         107 :             if (it->first.first == QEOpServerProxy::SUM) {
     102         159 :                 if (it->first.second == g_viz_constants.SESSION_FWD_TEARDOWN_BYTES ||
     103         150 :                     it->first.second == g_viz_constants.SESSION_FWD_TEARDOWN_PKTS ||
     104         141 :                     it->first.second == g_viz_constants.SESSION_REV_TEARDOWN_BYTES ||
     105         132 :                     it->first.second == g_viz_constants.SESSION_REV_TEARDOWN_PKTS ||
     106         225 :                     it->first.second == g_viz_constants.FLOW_TABLE_AGG_PKTS ||
     107          66 :                     it->first.second == g_viz_constants.FLOW_TABLE_AGG_BYTES) {
     108          18 :                     sname = it->first.second;
     109             :                 } else {
     110          66 :                     sname = string("SUM(") + it->first.second + string(")");
     111             :                 }
     112          23 :             } else if (it->first.first == QEOpServerProxy::COUNT) {
     113           7 :                 if (table == g_viz_constants.SESSION_SERIES_TABLE) {
     114           5 :                     sname = "sample_count";
     115             :                 } else {
     116           2 :                     sname = string("COUNT(") + it->first.second + string(")");
     117             :                 }
     118          16 :             } else if (it->first.first == QEOpServerProxy::CLASS) {
     119           9 :                 if (table == g_viz_constants.SESSION_SERIES_TABLE) {
     120           0 :                     sname = "session_class_id";
     121           9 :                 } else if (table == g_viz_constants.FLOW_SERIES_TABLE){
     122           0 :                     sname = "flow_class_id";
     123             :                 } else {
     124           9 :                     sname = string("CLASS(") + it->first.second + string(")");
     125             :                 }
     126           7 :             } else if (it->first.first == QEOpServerProxy::MAX) {
     127           1 :                 sname = string("MAX(") + it->first.second + string(")");
     128           6 :             } else if (it->first.first == QEOpServerProxy::MIN) {
     129           1 :                 sname = string("MIN(") + it->first.second + string(")");
     130           5 :             } else if (it->first.first == QEOpServerProxy::PERCENTILES) {
     131           1 :                 sname = string("PERCENTILES(") + it->first.second + string(")");
     132           4 :             } else if (it->first.first == QEOpServerProxy::AVG) {
     133           2 :                 sname = string("AVG(") + it->first.second + string(")");
     134           2 :             } else if (it->first.first == QEOpServerProxy::COUNT_DISTINCT) {
     135           2 :                 sname = string("COUNT_DISTINCT(") + it->first.second + string(")");
     136             :             } else {
     137           0 :                 QE_ASSERT(0);
     138             :             }
     139         107 :             switch (it->second.which()) {
     140             : 
     141         103 :                 case QEOpServerProxy::UINT64 : {
     142         103 :                         uint64_t mapit = boost::get<uint64_t>(it->second);
     143         103 :                         contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
     144         103 :                         val.SetUint64(mapit);
     145         103 :                         contrail_rapidjson::Value vk;
     146         103 :                         dd.AddMember(vk.SetString(sname.c_str(),
     147             :                                                   dd.GetAllocator()),
     148             :                                      val, dd.GetAllocator());
     149         103 :                     }
     150         103 :                     break;
     151           1 :                 case QEOpServerProxy::DOUBLE : {
     152           1 :                         double mapit = boost::get<double>(it->second);
     153           1 :                         contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
     154           1 :                         val.SetDouble(mapit);
     155           1 :                         contrail_rapidjson::Value vk;
     156           1 :                         dd.AddMember(vk.SetString(sname.c_str(),
     157             :                                                   dd.GetAllocator()),
     158             :                                      val, dd.GetAllocator());
     159           1 :                     }
     160           1 :                     break;
     161           1 :                 case QEOpServerProxy::TDIGEST : {
     162             :                         boost::shared_ptr<TDigest> mapit =
     163           1 :                             boost::get<boost::shared_ptr<TDigest> >(it->second);
     164             :                         
     165           1 :                         contrail_rapidjson::Value val(contrail_rapidjson::kObjectType);
     166           1 :                         float ptiles[7] = {.01,.05,.25,.5,.75,.95,.99};
     167           9 :                         std::string stiles[7] = {"01","05","25","50","75","95","99"};
     168             : #if 0
     169             :                         size_t jj = TDigest_get_ncentroids(mapit.get());
     170             :                         for (size_t kk=0; kk < jj; kk++) {
     171             :                             Centroid *c = TDigest_get_centroid(mapit.get(), kk);
     172             :                             std::cout << "Parse Centroid " <<
     173             :                                Centroid_get_mean(c) << " " <<
     174             :                                Centroid_get_count(c) << " " <<
     175             :                                Centroid_quantile(c, mapit.get()) << std::endl;
     176             :                         }
     177             : #endif
     178           8 :                         for (size_t idx=0; idx<7; idx++) {
     179           7 :                             contrail_rapidjson::Value sval(contrail_rapidjson::kNumberType);
     180           7 :                             sval.SetDouble(TDigest_percentile(mapit.get(), ptiles[idx]));
     181           7 :                             contrail_rapidjson::Value vk;
     182           7 :                             val.AddMember(vk.SetString(stiles[idx].c_str(),
     183             :                                                        dd.GetAllocator()),
     184             :                                           sval, dd.GetAllocator());
     185           7 :                         }
     186           1 :                         contrail_rapidjson::Value vk;
     187           1 :                         dd.AddMember(vk.SetString(sname.c_str(),
     188             :                                                   dd.GetAllocator()),
     189             :                                      val, dd.GetAllocator());
     190          10 :                     }
     191           1 :                     break;
     192           2 :                 case QEOpServerProxy::CENTROID : {
     193             :                         boost::shared_ptr<Centroid> mapit =
     194           2 :                             boost::get<boost::shared_ptr<Centroid> >(it->second);
     195           2 :                         contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
     196           2 :                         val.SetDouble(Centroid_get_mean(mapit.get()));
     197           2 :                         contrail_rapidjson::Value vk;
     198           2 :                         dd.AddMember(vk.SetString(sname.c_str(),
     199             :                                                   dd.GetAllocator()),
     200             :                                      val, dd.GetAllocator());
     201           2 :                     }
     202           2 :                     break;
     203           0 :                 default:
     204           0 :                     QE_ASSERT(0);
     205             :             }
     206         107 :         }
     207           0 :     } catch (boost::bad_get& ex) {
     208           0 :         QE_ASSERT(0);
     209           0 :     } catch (const std::out_of_range& oor) {
     210           0 :         QE_ASSERT(0);
     211             :     }
     212         392 :     contrail_rapidjson::StringBuffer sb;
     213         196 :     contrail_rapidjson::Writer<contrail_rapidjson::StringBuffer> writer(sb);
     214         196 :     dd.Accept(writer);
     215         196 :     jstr = sb.GetString();
     216             :     //QE_LOG_NOQID(INFO, "Jsonify row " << jstr);
     217         392 :     return true;
     218         196 : }
     219             : 
     220          64 : void StatsSelect::MergeFinal(const std::vector<boost::shared_ptr<MapBufT> >& inputs,
     221             :         MapBufT& output) {
     222          64 :     MapBufT temp;
     223         576 :     for (size_t idx = 0; idx < inputs.size(); idx++) {
     224         512 :         if (agg_sort_cols_.empty())
     225         504 :             Merge(count_distinct_field_, *inputs[idx], output);
     226             :         else
     227           8 :             Merge(count_distinct_field_, *inputs[idx], temp);
     228             :     }
     229          64 :     if (!agg_sort_cols_.empty()) {
     230           1 :         MapBufT::iterator jt = temp.end();
     231           4 :         for (MapBufT::iterator it = temp.begin(); it!= temp.end(); it++) {
     232           3 :             if (jt!=temp.end()) {
     233           2 :                 temp.erase(jt);
     234             :             }
     235           3 :             std::vector<StatVal> ukey = it->first;
     236           3 :             const StatMap& uniks = it->second.first;
     237           3 :             const QEOpServerProxy::AggRowT& narows = it->second.second;
     238           3 :             for (AggSortT::const_iterator xt = agg_sort_cols_.begin();
     239           6 :                     xt != agg_sort_cols_.end(); xt++) {
     240           3 :                 QE_ASSERT(narows.find(xt->first) != narows.end());
     241           3 :                 ukey[xt->second] = narows.at(xt->first);
     242             :             }
     243           3 :             MergeFullRow(ukey, uniks, narows, output);
     244           3 :             jt = it;
     245           3 :         }
     246           1 :         temp.clear();
     247             :     }
     248          64 : }
     249             : 
     250         936 : void StatsSelect::Merge(const std::string& count_distinct_field_, const MapBufT& input, MapBufT& output) {
     251         936 :     for (MapBufT::const_iterator it = input.begin();
     252        1490 :             it!=input.end(); it++) {
     253         554 :         std::vector<StatVal> ukey(it->first);
     254         554 :         StatMap uniks(it->second.first);
     255         554 :         QEOpServerProxy::AggRowT narows(it->second.second);
     256             : 
     257         554 :         StatMap::iterator itr = uniks.find(count_distinct_field_);
     258         554 :         if (itr != uniks.end()) {
     259             :             pair<QEOpServerProxy::AggOper, std::string> aggkey(
     260           3 :                 QEOpServerProxy::COUNT_DISTINCT, count_distinct_field_);
     261           3 :             narows.insert(make_pair(aggkey, (uint64_t) 1));
     262           3 :             uniks.erase(itr);
     263           3 :             size_t hash_slot = ukey.size() - 1;
     264           3 :             uint64_t hash_val = boost::hash_range(uniks.begin(), uniks.end());
     265           3 :             ukey[hash_slot] = hash_val;
     266           3 :         }
     267             : 
     268         554 :         MergeFullRow(ukey, uniks, narows, output);
     269             : 
     270         554 :     }
     271         936 : }
     272             : 
     273        1073 : StatsSelect::StatsSelect(AnalyticsQuery * m_query,
     274        1073 :                 const std::vector<std::string> & select_fields) :
     275        1073 :                         main_query(m_query), select_fields_(select_fields),
     276        1074 :                         ts_period_(0), isT_(false), count_field_() {
     277             : 
     278        1073 :     status_ = false;
     279             : 
     280        5197 :     for (size_t j=0; j<select_fields_.size(); j++) {
     281             : 
     282        4122 :         if (select_fields_[j] == g_viz_constants.STAT_TIME_FIELD) {
     283          89 :             if (ts_period_) {
     284           0 :                 QE_LOG(INFO, "StatsSelect cannot accept both T and T="); 
     285           0 :                 return;
     286             :             }
     287          89 :             isT_ = true;
     288          89 :             sort_cols_.insert(std::make_pair(select_fields_[j], 0));
     289          89 :             QE_TRACE(DEBUG, "StatsSelect T");
     290             : 
     291        4035 :         } else if (select_fields_[j].substr(0, g_viz_constants.STAT_TIMEBIN_FIELD.size()) ==
     292             :                 g_viz_constants.STAT_TIMEBIN_FIELD) {
     293          40 :             if (isT_) {
     294           0 :                 QE_LOG(INFO, "StatsSelect cannot accept both T and T="); 
     295           0 :                 return;
     296             :             }
     297          40 :             string tsstr = select_fields_[j].substr(g_viz_constants.STAT_TIMEBIN_FIELD.size());
     298          40 :             ts_period_ = strtoul(tsstr.c_str(), NULL, 10) * 1000000;
     299          40 :             if (!ts_period_) {
     300           0 :                 QE_LOG(INFO, "StatsSelect cannot accept T= for " << tsstr);
     301           0 :                 return;
     302             :             }
     303          40 :             sort_cols_.insert(std::make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,0));
     304          40 :             QE_TRACE(DEBUG, "StatsSelect T=");
     305          40 :         } else {
     306        3995 :             std::string sfield = select_fields_[j];
     307             :             QEOpServerProxy::AggOper agg =
     308        3995 :                 StatsQuery::ParseAgg(select_fields_[j], sfield);
     309        3995 :             if (agg == QEOpServerProxy::INVALID) {
     310        3474 :                 unik_cols_.insert(select_fields_[j]);
     311        3473 :                 QE_TRACE(DEBUG, "StatsSelect unik field " << select_fields_[j]);
     312         521 :             } else if (agg == QEOpServerProxy::COUNT) {
     313          54 :                 count_field_ = sfield;
     314          54 :                 QE_TRACE(DEBUG, "StatsSelect COUNT " << sfield);
     315         467 :             } else if (agg == QEOpServerProxy::COUNT_DISTINCT) {
     316          18 :                 count_distinct_field_ = sfield;
     317          18 :                 unik_cols_.insert(sfield);
     318          18 :                 QE_TRACE(DEBUG, "StatsSelect COUNT_DISTINCT " << sfield);
     319         449 :             } else if (agg == QEOpServerProxy::SUM) {
     320         288 :                 sum_cols_.insert(sfield);
     321         288 :                 QE_TRACE(DEBUG, "StatsSelect SUM " << sfield);
     322         161 :             } else if (agg == QEOpServerProxy::CLASS) {
     323          17 :                 class_cols_.insert(sfield);
     324          17 :                 QE_TRACE(DEBUG, "StatsSelect CLASS " << sfield);
     325         144 :             } else if (agg == QEOpServerProxy::MAX) {
     326          36 :                 max_field_.insert(sfield);
     327          36 :                 QE_TRACE(DEBUG, "StatsSelect MAX " << sfield);
     328         108 :             } else if (agg == QEOpServerProxy::MIN) {
     329          18 :                 min_field_.insert(sfield);
     330          18 :                 QE_TRACE(DEBUG, "StatsSelect MIN " << sfield);
     331          90 :             } else if (agg == QEOpServerProxy::PERCENTILES) {
     332          36 :                 percentile_cols_.insert(sfield);
     333          36 :                 QE_TRACE(DEBUG, "StatsSelect PERCENTILES " << sfield);
     334          54 :             } else if (agg == QEOpServerProxy::AVG) {
     335          54 :                 avg_field_.insert(sfield);
     336          54 :                 QE_TRACE(DEBUG, "StatsSelect AVG " << sfield);
     337             :             } else {
     338           0 :                 QE_ASSERT(0);
     339             :             }
     340        3995 :         }
     341             :     }
     342             : 
     343        1075 :     status_ = true;
     344           0 : }
     345             : 
     346        1075 : void StatsSelect::SetSortOrder(const std::vector<sort_field_t>& sort_fields) {
     347        1075 :     if (sort_fields.size()) {
     348          24 :         sort_cols_.clear();
     349             :     } else {
     350        1051 :         QE_TRACE(DEBUG, "StatsSelect Sort By T/T=");
     351             :     }
     352        1099 :     for (size_t st = 0; st < sort_fields.size(); st++) {
     353             : 
     354          24 :         if (sort_fields[st].name.substr(0, g_viz_constants.STAT_TIMEBIN_FIELD.size()) ==
     355             :                 g_viz_constants.STAT_TIMEBIN_FIELD) {
     356           0 :             QE_TRACE(DEBUG, "StatsSelect Sort By T=");
     357           0 :             sort_cols_.insert(std::make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,st));
     358          24 :         } else if (sort_fields[st].name == g_viz_constants.STAT_TIME_FIELD) {
     359           0 :             QE_TRACE(DEBUG, "StatsSelect Sort By " << sort_fields[st].name);
     360           0 :             sort_cols_.insert(std::make_pair(sort_fields[st].name, st));                
     361             :         } else {
     362          24 :             std::string sfield = sort_fields[st].name;
     363             :             QEOpServerProxy::AggOper agg =
     364          24 :                 StatsQuery::ParseAgg(sort_fields[st].name, sfield);
     365          24 :             if (main_query->is_stat_table_query(main_query->table()) &&
     366           0 :                     main_query->stats().is_stat_table_static()) {
     367           0 :                 StatsQuery::column_t c = main_query->stats().get_column_desc(sfield);
     368           0 :                 if (c.datatype == QEOpServerProxy::BLANK) {
     369           0 :                     QE_LOG(INFO, "StatsSelect unknown field " << sort_fields[st].name);
     370           0 :                     return;
     371             :                 }
     372           0 :             }
     373          24 :             if (agg == QEOpServerProxy::INVALID) {
     374          12 :                 QE_TRACE(DEBUG, "StatsSelect Sort By " << sort_fields[st].name);
     375          12 :                 sort_cols_.insert(std::make_pair(sort_fields[st].name,st));                
     376             :             } else  {
     377          12 :                 agg_sort_cols_.insert(make_pair(make_pair(agg, sfield), st));
     378          12 :                 QE_TRACE(DEBUG, "StatsSelect Sort Agg " << agg << " of " << sfield);
     379             :             }
     380          24 :         }
     381             :     }    
     382             : }
     383             : 
     384        1092 : void StatsSelect::MergeFullRow(
     385             :         const std::vector<StatVal>& ukey,
     386             :         const StatMap& uniks,
     387             :         const QEOpServerProxy::AggRowT& narows,
     388             :         MapBufT& output) {
     389             : 
     390        1092 :     MapBufT::iterator rt = output.find(ukey);
     391        1092 :     if (rt!= output.end()) {
     392         227 :         StatMap& rm = rt->second.first;
     393         227 :         if (rm!=uniks) {
     394           0 :             output.insert(make_pair(ukey, make_pair(uniks, narows)));
     395             :         } else {
     396         226 :             QEOpServerProxy::AggRowT &arows = rt->second.second;
     397         226 :             MergeAggRow(arows, narows);
     398             :         }
     399             :     } else {
     400         865 :         QEOpServerProxy::AggRowT arows;
     401         865 :         output.insert(make_pair(ukey, make_pair(uniks, narows)));
     402         865 :     }    
     403        1091 : }
     404             : 
     405             : 
     406         226 : void StatsSelect::MergeAggRow(QEOpServerProxy::AggRowT &arows,
     407             :         const QEOpServerProxy::AggRowT &narows) {
     408         226 :     if (arows.empty()) {
     409         108 :         arows = narows;
     410         108 :         return;
     411             :     }
     412         118 :     for (QEOpServerProxy::AggRowT::iterator jt = arows.begin();
     413         361 :             jt!= arows.end(); jt++) {
     414         243 :         QEOpServerProxy::AggRowT::const_iterator kt = narows.find(jt->first);
     415         244 :         if (kt!=narows.end()) {
     416             :             // Attribute name must match for aggregate and for the new value
     417         208 :             QE_ASSERT(jt->first.second == kt->first.second);
     418             : 
     419         207 :             if (jt->first.first == QEOpServerProxy::SUM) {
     420         173 :                 StatsSelect::StatVal & sv = jt->second;
     421             :                 try {
     422         173 :                     if (sv.which() == QEOpServerProxy::UINT64) {
     423         173 :                         sv = boost::get<uint64_t>(sv) + boost::get<uint64_t>(kt->second);
     424             :                     }
     425         174 :                     if (sv.which() == QEOpServerProxy::DOUBLE) {
     426           0 :                         sv = boost::get<double>(sv) + boost::get<double>(kt->second);
     427             :                     }                   
     428           0 :                 } catch (boost::bad_get& ex) {
     429           0 :                     QE_ASSERT(0);
     430           0 :                 } catch (const std::out_of_range& oor) {
     431           0 :                     QE_ASSERT(0);
     432             :                 }                        
     433             : 
     434             :             }
     435             : 
     436         392 :             if (jt->first.first == QEOpServerProxy::COUNT ||
     437         185 :                 jt->first.first == QEOpServerProxy::COUNT_DISTINCT) {
     438             :                 try {
     439          24 :                     uint64_t& sv =  boost::get<uint64_t>(jt->second);
     440          24 :                     sv = sv + boost::get<uint64_t>(kt->second);
     441           0 :                 } catch (boost::bad_get& ex) {
     442           0 :                     QE_ASSERT(0);
     443           0 :                 } catch (const std::out_of_range& oor) {
     444           0 :                     QE_ASSERT(0);
     445             :                 }     
     446             :             }
     447         207 :             if (jt->first.first == QEOpServerProxy::MAX) {
     448           2 :                 StatsSelect::StatVal & sv = jt->second;
     449             :                 try {
     450           2 :                     if (sv.which() == QEOpServerProxy::UINT64) {
     451           2 :                         uint64_t& existing_max = boost::get<uint64_t>(jt->second);
     452           2 :                         existing_max = (existing_max > (boost::get<uint64_t>(kt->second))) ? existing_max:(boost::get<uint64_t>(kt->second));
     453             :                     }
     454           2 :                     if (sv.which() == QEOpServerProxy::DOUBLE) {
     455           0 :                         double& existing_max = boost::get<double>(jt->second);
     456           0 :                         existing_max = (existing_max > (boost::get<double>(kt->second))) ? existing_max:(boost::get<double>(kt->second));
     457             :                     }                   
     458           0 :                 } catch (boost::bad_get& ex) {
     459           0 :                     QE_ASSERT(0);
     460             :                 }     
     461             :             }
     462         207 :             if (jt->first.first == QEOpServerProxy::MIN) {
     463           2 :                 StatsSelect::StatVal & sv = jt->second;
     464             :                 try {
     465           2 :                     if (sv.which() == QEOpServerProxy::UINT64) {
     466           0 :                         uint64_t& existing_min = boost::get<uint64_t>(jt->second);
     467           0 :                         existing_min = (existing_min < (boost::get<uint64_t>(kt->second))) ? existing_min:(boost::get<uint64_t>(kt->second));
     468             :                     }
     469           2 :                     if (sv.which() == QEOpServerProxy::DOUBLE) {
     470           2 :                         double& existing_min = boost::get<double>(jt->second);
     471           2 :                         existing_min = (existing_min < (boost::get<double>(kt->second))) ? existing_min:(boost::get<double>(kt->second));
     472             :                     }                   
     473           0 :                 } catch (boost::bad_get& ex) {
     474           0 :                     QE_ASSERT(0);
     475             :                 }     
     476             :             }
     477         207 :             if (jt->first.first == QEOpServerProxy::AVG) {
     478           4 :                 StatsSelect::StatVal & sv = jt->second;
     479             :                 try {
     480           4 :                     if (sv.which() == QEOpServerProxy::CENTROID) {
     481             :                         boost::shared_ptr<Centroid>& pt =
     482           4 :                             boost::get<boost::shared_ptr<Centroid> >(jt->second);
     483             : 
     484             :                         boost::shared_ptr<Centroid> pt2 =
     485           4 :                             boost::get<boost::shared_ptr<Centroid> >(kt->second);
     486           4 :                         Centroid_add(pt.get(),
     487             :                                 Centroid_get_mean(pt2.get()),
     488             :                                 Centroid_get_count(pt2.get()));
     489           4 :                     }
     490           0 :                 } catch (boost::bad_get& ex) {
     491           0 :                     QE_ASSERT(0);
     492             :                 }
     493             :             }
     494         207 :             if (jt->first.first == QEOpServerProxy::PERCENTILES) {
     495           2 :                 StatsSelect::StatVal & sv = jt->second;
     496             :                 try {
     497           2 :                     if (sv.which() == QEOpServerProxy::TDIGEST) {
     498             :                         boost::shared_ptr<TDigest>& pt =
     499           2 :                             boost::get<boost::shared_ptr<TDigest> >(jt->second);
     500             :                         
     501             :                         boost::shared_ptr<TDigest> pt2 =
     502           2 :                             boost::get<boost::shared_ptr<TDigest> >(kt->second);
     503           2 :                         size_t j, ncentroids = TDigest_get_ncentroids(pt2.get());
     504           4 :                         for (j = 0; j < ncentroids; j++) {
     505           2 :                             Centroid * c = TDigest_get_centroid(pt2.get(), j);
     506           2 :                             TDigest* nd = TDigest_add(pt.get(),
     507             :                                 Centroid_get_mean(c),
     508             :                                 Centroid_get_count(c));
     509           2 :                             if (nd) {
     510           0 :                                 pt.reset(nd);
     511             :                             }
     512             :                         }
     513           2 :                     }
     514           0 :                 } catch (boost::bad_get& ex) {
     515           0 :                     QE_ASSERT(0);
     516             :                 }     
     517             :             }
     518             :         }
     519             :     }
     520             : }
     521             : 
     522             : namespace boost {
     523             :    std::size_t hash_value(const StatsSelect::StatVal&); 
     524             : }
     525             : 
     526        1822 : std::size_t boost::hash_value(const StatsSelect::StatVal& sv) {
     527        1822 :     std::ostringstream ostr;
     528        1822 :     ostr << sv;
     529        3643 :     return boost::hash_value(ostr.str());
     530        1821 : }
     531             : 
     532         535 : bool StatsSelect::LoadRow(boost::uuids::uuid u,
     533             :                 uint64_t timestamp, const vector<StatEntry>& row, MapBufT& output) {
     534             : 
     535         535 :         if (!Status()) return false;
     536         535 :     uint64_t ts = 0;
     537             : 
     538             :     // Build Uniks map
     539         535 :     StatMap uniks;
     540             :     set<string>::const_iterator ukit =
     541         535 :         unik_cols_.find(g_viz_constants.STAT_UUID_FIELD);
     542         535 :     if (ukit!=unik_cols_.end()) {
     543          63 :         uniks.insert(make_pair(g_viz_constants.STAT_UUID_FIELD,u));
     544             :     }
     545             :     
     546         535 :     if (isT_) {
     547          50 :         ts = timestamp;
     548          50 :         uniks.insert(make_pair(g_viz_constants.STAT_TIME_FIELD,ts)); 
     549             :     }
     550             : 
     551         535 :     if (ts_period_) {
     552          54 :         ts = timestamp - (timestamp % ts_period_);
     553          54 :         uniks.insert(make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,ts)); 
     554             :     }
     555             : 
     556         535 :     for (vector<StatEntry>::const_iterator it = row.begin();
     557       18104 :             it != row.end(); it++) {
     558       17571 :         set<string>::const_iterator uit = unik_cols_.find(it->name);
     559       17568 :         if (uit!=unik_cols_.end()) {
     560        1633 :             uniks.insert(make_pair(it->name, it->value));
     561             :         }
     562             :     }
     563             : 
     564             :     // Build sort vector
     565             :     // Last slot is reserved for the hash
     566         535 :     std::vector<StatVal> ukey(sort_cols_.size() + agg_sort_cols_.size() + 1);
     567         534 :     size_t hash_slot = sort_cols_.size() + agg_sort_cols_.size();
     568         534 :     uint64_t hash_val = boost::hash_range(uniks.begin(), uniks.end());
     569         535 :     ukey[hash_slot] = hash_val;
     570             : 
     571         534 :     for (map<string, size_t>::const_iterator st = sort_cols_.begin();
     572         664 :             st!=sort_cols_.end(); st++) {
     573         130 :         QE_ASSERT(uniks.find(st->first) != uniks.end());
     574         130 :         ukey[st->second] = uniks.at(st->first);
     575             :     }
     576             : 
     577         534 :     QEOpServerProxy::AggRowT narows;
     578         534 :     for (vector<StatEntry>::const_iterator it = row.begin();
     579       18115 :             it != row.end(); it++) {
     580       17577 :         set<string>::const_iterator uit = sum_cols_.find(it->name);
     581       17580 :         if (uit!=sum_cols_.end()) {
     582         270 :             pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::SUM,it->name);
     583         270 :             narows.insert(make_pair(aggkey, it->value)); 
     584         271 :         }
     585             :     }
     586             : 
     587         535 :     for (vector<StatEntry>::const_iterator it = row.begin();
     588       18129 :             it != row.end(); it++) {
     589       17594 :         set<string>::const_iterator uit = max_field_.find(it->name);
     590       17593 :         if (uit!=max_field_.end()) {
     591           3 :             pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::MAX,it->name);
     592           3 :             narows.insert(make_pair(aggkey, it->value)); 
     593           3 :         }
     594       17590 :         uit = min_field_.find(it->name);
     595       17589 :         if (uit!=min_field_.end()) {
     596           3 :             pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::MIN,it->name);
     597           3 :             narows.insert(make_pair(aggkey, it->value)); 
     598           3 :         }
     599       17589 :         uit = avg_field_.find(it->name);
     600       17595 :         if (uit!=avg_field_.end()) {
     601           6 :             pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::AVG,it->name);
     602             :             try {
     603           6 :                 Centroid *t = Centroid_create(0, 0);
     604           6 :                 StatsSelect::StatVal sv = it->value;
     605             :                 double val;
     606           6 :                 if (sv.which() == QEOpServerProxy::UINT64) {
     607           3 :                     val = boost::get<uint64_t>(sv);
     608             :                 } else {
     609           3 :                     val = boost::get<double>(sv);
     610             :                 }
     611           6 :                 Centroid_add(t, val, 1);
     612           6 :                 boost::shared_ptr<Centroid> pt(t, &StatsSelect::DeleteCentroid);
     613           6 :                 narows.insert(make_pair(aggkey, pt));
     614           6 :             } catch (boost::bad_get& ex) {
     615           0 :                 QE_ASSERT(0);
     616           0 :             } catch (const std::out_of_range& oor) {
     617           0 :                 QE_ASSERT(0);
     618             :             }
     619           6 :         }
     620       17594 :         uit = percentile_cols_.find(it->name);
     621       17594 :         if (uit!=percentile_cols_.end()) {
     622           3 :             pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::PERCENTILES,it->name);
     623             : 
     624             :             try {
     625           3 :                 TDigest *t = TDigest_create(0.01, 100);
     626           3 :                 StatsSelect::StatVal sv = it->value;
     627             :                 double val;
     628           3 :                 if (sv.which() == QEOpServerProxy::UINT64) {
     629           3 :                     val = boost::get<uint64_t>(sv);
     630             :                 } else {
     631           0 :                     val = boost::get<double>(sv);
     632             :                 }                   
     633           3 :                 TDigest_add(t, val, 1);
     634           3 :                 boost::shared_ptr<TDigest> pt(t, &StatsSelect::DeleteTDigest);
     635           3 :                 narows.insert(make_pair(aggkey, pt)); 
     636           3 :             } catch (boost::bad_get& ex) {
     637           0 :                 QE_ASSERT(0);
     638           0 :             } catch (const std::out_of_range& oor) {
     639           0 :                 QE_ASSERT(0);
     640             :             }                        
     641           3 :         }
     642             :     }
     643             :     
     644         535 :     for (std::set<std::string>::const_iterator ct = class_cols_.begin();
     645         544 :             ct!=class_cols_.end(); ct++) {
     646           9 :         pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::CLASS,*ct);
     647           9 :         StatMap huniks;
     648           9 :         for (vector<StatEntry>::const_iterator rit = row.begin();
     649          90 :                 rit != row.end(); rit++) {
     650          81 :             if (rit->name != *ct) {
     651          81 :                 if (uniks.find(rit->name) != uniks.end()) {
     652             :                     // For generating the hash, consider all attributes that 
     653             :                     // are in the row, and that do not match the CLASS attribute,
     654             :                     // and that are in non-aggregate attributes in the SELECT
     655          18 :                     huniks[rit->name] = rit->value;
     656             :                 }
     657             :             }
     658             :         }
     659           9 :         uint64_t hh = boost::hash_range(huniks.begin(), huniks.end());
     660           9 :         narows.insert(make_pair(aggkey, hh));
     661           9 :     }
     662             :  
     663         535 :     if (!count_field_.empty()) {
     664          30 :         pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::COUNT,count_field_);
     665          30 :         narows.insert(make_pair(aggkey, (uint64_t) 1));            
     666          30 :     }
     667             :     
     668         535 :     MergeFullRow(ukey, uniks, narows, output);
     669             : 
     670         534 :     return true;
     671         534 : }
     672             : 

Generated by: LCOV version 1.14