LCOV - code coverage report
Current view: top level - xmpp - xmpp_server.cc (source / functions) Hit Total Coverage
Test: OpenSDN C/C++ coverage (all TARGET_SET jobs) Lines: 326 499 65.3 %
Date: 2026-10-05 02:12:29 Functions: 48 66 72.7 %
Legend: Lines: hit not hit

          Line data    Source code
       1             : /*
       2             :  * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
       3             :  */
       4             : 
       5             : #include "xmpp/xmpp_server.h"
       6             : 
       7             : #include <boost/foreach.hpp>
       8             : #include <boost/tuple/tuple.hpp>
       9             : 
      10             : #include "base/task_annotations.h"
      11             : #include "xmpp/xmpp_connection.h"
      12             : #include "xmpp/xmpp_factory.h"
      13             : #include "xmpp/xmpp_lifetime.h"
      14             : #include "xmpp/xmpp_log.h"
      15             : #include "xmpp/xmpp_sandesh.h"
      16             : #include "xmpp/xmpp_session.h"
      17             : 
      18             : #include "sandesh/request_pipeline.h"
      19             : #include "sandesh/common/vns_types.h"
      20             : #include "sandesh/common/vns_constants.h"
      21             : #include "sandesh/xmpp_server_types.h"
      22             : #include "sandesh/xmpp_trace_sandesh_types.h"
      23             : #include "sandesh/xmpp_client_server_sandesh_types.h"
      24             : 
      25             : using namespace std;
      26             : using namespace boost::asio;
      27             : using boost::tie;
      28             : 
      29             : #define DEFAULT_XMPP_HOLD_TIME 90
      30             : class XmppServer::DeleteActor : public LifetimeActor {
      31             : public:
      32         120 :     DeleteActor(XmppServer *server)
      33         120 :         : LifetimeActor(server->lifetime_manager()), server_(server) {
      34         120 :     }
      35         118 :     virtual bool MayDelete() const {
      36         118 :         CHECK_CONCURRENCY("bgp::Config");
      37         118 :         return server_->MayDelete();
      38             :     }
      39         118 :     virtual void Shutdown() {
      40         118 :         CHECK_CONCURRENCY("bgp::Config");
      41         118 :         server_->SessionShutdown();
      42         118 :     }
      43         118 :     virtual void Destroy() {
      44         118 :         CHECK_CONCURRENCY("bgp::Config");
      45         118 :         server_->Terminate();
      46         118 :     }
      47             : 
      48             : private:
      49             :     XmppServer *server_;
      50             : };
      51             : 
      52             : 
      53          74 : XmppServer::XmppServer(EventManager *evm, const string &server_addr,
      54          74 :                        const XmppChannelConfig *config)
      55             :     : XmppConnectionManager(
      56          74 :           evm, ssl::context::sslv23_server, config->auth_enabled, true),
      57          74 :       max_connections_(0),
      58          74 :       lifetime_manager_(XmppStaticObjectFactory::Create<XmppLifetimeManager>(
      59             :           TaskScheduler::GetInstance()->GetTaskId("bgp::Config"))),
      60          74 :       deleter_(new DeleteActor(this)),
      61          74 :       server_addr_(server_addr),
      62          74 :       log_uve_(false),
      63          74 :       auth_enabled_(config->auth_enabled),
      64          74 :       tcp_hold_time_(config->tcp_hold_time),
      65          74 :       gr_helper_disable_(config->gr_helper_disable),
      66          74 :       dscp_value_(0),
      67          74 :       connection_queue_(TaskScheduler::GetInstance()->GetTaskId("bgp::Config"),
      68         222 :           0, boost::bind(&XmppServer::DequeueConnection, this, _1)) {
      69             : 
      70          74 :     if (config->auth_enabled) {
      71             : 
      72             :         // Get SSL context from base class and update
      73          54 :         boost::asio::ssl::context *ctx = context();
      74          54 :         boost::system::error_code ec;
      75             : 
      76             :         // set mode
      77          54 :         ctx->set_options(ssl::context::default_workarounds |
      78             :                          ssl::context::no_sslv3 | ssl::context::no_sslv2, ec);
      79          54 :         if (ec.value() != 0) {
      80           0 :             LOG(ERROR, "Error : " << ec.message() << ", setting ssl options");
      81           0 :             exit(EINVAL);
      82             :         }
      83             : 
      84             :         // CA certificate, used to verify if the peer certificate
      85             :         // is signed by a trusted CA
      86          54 :         std::string ca_cert_filename = config->path_to_ca_cert;
      87          54 :         if (!ca_cert_filename.empty()) {
      88             : 
      89             :             // Verify peer has CA signed certificate
      90           0 :             ctx->set_verify_mode(boost::asio::ssl::verify_peer, ec);
      91           0 :             if (ec.value() != 0) {
      92           0 :                 LOG(ERROR, "Error : " << ec.message()
      93             :                     << ", while setting ssl verification mode");
      94           0 :                 exit(EINVAL);
      95             :             }
      96             : 
      97           0 :             ctx->load_verify_file(config->path_to_ca_cert, ec);
      98           0 :             if (ec.value() != 0) {
      99           0 :                 LOG(ERROR, "Error : " << ec.message()
     100             :                     << ", while using cacert file : "
     101             :                     << config->path_to_ca_cert);
     102           0 :                 exit(EINVAL);
     103             :             }
     104             :         }
     105             : 
     106             :         // server certificate
     107          54 :         ctx->use_certificate_file(config->path_to_server_cert,
     108             :                                   boost::asio::ssl::context::pem, ec);
     109          54 :         if (ec.value() != 0) {
     110           0 :             LOG(ERROR, "Error : " << ec.message()
     111             :                 << ", while using server cert file : "
     112             :                 << config->path_to_server_cert);
     113           0 :             exit(EINVAL);
     114             :         }
     115             : 
     116             :         // server private key
     117          54 :         ctx->use_private_key_file(config->path_to_server_priv_key,
     118             :                                   boost::asio::ssl::context::pem, ec);
     119          54 :         if (ec.value() != 0) {
     120           0 :             LOG(ERROR, "Error : " << ec.message()
     121             :                 << ", while using privkey file : "
     122             :                 << config->path_to_server_priv_key);
     123           0 :             exit(EINVAL);
     124             :         }
     125          54 :     }
     126          74 : }
     127             : 
     128             : class XmppConfigUpdater {
     129             : public:
     130          40 :     explicit XmppConfigUpdater(XmppServer *server,
     131          40 :                                BgpConfigManager *config_manager) :
     132          40 :             server_(server) {
     133          40 :         BgpConfigManager::Observers obs;
     134             :         obs.system= boost::bind(&XmppConfigUpdater::ProcessGlobalSystemConfig,
     135          40 :             this, _1, _2);
     136             :         obs.protocol = boost::bind(&XmppConfigUpdater::ProcessProtocolConfig,
     137          40 :             this, _1, _2);
     138          40 :         config_manager->RegisterObservers(obs);
     139          40 :     }
     140             : 
     141        2300 :     const BgpGlobalSystemConfig &config() const { return config_; }
     142             : 
     143          40 :     void ProcessProtocolConfig(const BgpProtocolConfig *protocol_config,
     144             :                                BgpConfigManager::EventType event) {
     145             : 
     146          40 :         if (server_->subcluster_name() != protocol_config->subcluster_name()) {
     147           0 :             server_->set_subcluster_name(protocol_config->subcluster_name());
     148           0 :             server_->ClearAllConnections();
     149             :         }
     150          40 :     }
     151             : 
     152          60 :     void ProcessGlobalSystemConfig(const BgpGlobalSystemConfig *system,
     153             :             BgpConfigManager::EventType event) {
     154             :         // Clear peers only if GR is or was enabled.
     155          60 :         bool clear_peers = config_.gr_enable() || system->gr_enable();
     156          60 :         bool update_peers = false;
     157          60 :         config_.set_gr_enable(system->gr_enable());
     158          60 :         config_.set_gr_time(system->gr_time());
     159          60 :         config_.set_llgr_time(system->llgr_time());
     160          60 :         config_.set_end_of_rib_timeout(system->end_of_rib_timeout());
     161          60 :         config_.set_gr_xmpp_helper(system->gr_xmpp_helper());
     162             : 
     163             :         // Process any change in xmpp-hold-time
     164         120 :         if (config_.fc_enabled() != system->fc_enabled() ||
     165          60 :             config_.xmpp_hold_time() != system->xmpp_hold_time()) {
     166           0 :             config_.set_xmpp_hold_time(system->xmpp_hold_time());
     167           0 :             config_.set_fc_enabled(system->fc_enabled());
     168             :             // if fc_enabled is not set, use default hold time
     169           0 :             if (!system->fc_enabled())
     170           0 :                 config_.set_xmpp_hold_time(DEFAULT_XMPP_HOLD_TIME);
     171           0 :             update_peers = true;
     172             :         }
     173             : 
     174             :         // Process any change in rd-cluster-seed knob
     175          60 :         if (config_.rd_cluster_seed() != system->rd_cluster_seed()) {
     176           0 :             config_.set_rd_cluster_seed(system->rd_cluster_seed());
     177           0 :             clear_peers = true;
     178             :         }
     179             : 
     180          60 :         if (clear_peers)
     181          60 :             server_->ClearAllConnections();
     182           0 :         else if (update_peers)
     183           0 :             server_->UpdateAllConnections(config_.xmpp_hold_time());
     184          60 :     }
     185             : 
     186             :     const string subcluster_name() const { return server_->subcluster_name(); }
     187             :     void set_subcluster_name(const string& name) {
     188             :         server_->set_subcluster_name(name);
     189             :     }
     190             : 
     191         265 :     uint8_t xmpp_hold_time() const { return config_.xmpp_hold_time(); }
     192             :     void set_xmpp_hold_time(int hold_time) {
     193             :         config_.set_xmpp_hold_time(hold_time);
     194             :     }
     195             : 
     196             : private:
     197             :     XmppServer *server_;
     198             :     BgpGlobalSystemConfig config_;
     199             : };
     200             : 
     201          39 : XmppServer::XmppServer(EventManager *evm, const string &server_addr)
     202             :     : XmppConnectionManager(evm, ssl::context::sslv23_server, false, false),
     203          39 :       max_connections_(0),
     204          39 :       lifetime_manager_(XmppStaticObjectFactory::Create<XmppLifetimeManager>(
     205             :           TaskScheduler::GetInstance()->GetTaskId("bgp::Config"))),
     206          39 :       deleter_(new DeleteActor(this)),
     207          39 :       server_addr_(server_addr),
     208          39 :       log_uve_(false),
     209          39 :       auth_enabled_(false),
     210          39 :       tcp_hold_time_(XmppChannelConfig::kTcpHoldTime),
     211          39 :       gr_helper_disable_(false),
     212          39 :       xmpp_config_updater_(NULL),
     213          39 :       dscp_value_(0),
     214          39 :       connection_queue_(TaskScheduler::GetInstance()->GetTaskId("bgp::Config"),
     215         117 :           0, boost::bind(&XmppServer::DequeueConnection, this, _1)) {
     216          39 : }
     217             : 
     218             : 
     219           7 : XmppServer::XmppServer(EventManager *evm)
     220             :     : XmppConnectionManager(evm, ssl::context::sslv23_server, false, false),
     221           7 :       max_connections_(0),
     222           7 :       lifetime_manager_(new LifetimeManager(
     223          14 :           TaskScheduler::GetInstance()->GetTaskId("bgp::Config"))),
     224           7 :       deleter_(new DeleteActor(this)),
     225           7 :       log_uve_(false),
     226           7 :       auth_enabled_(false),
     227           7 :       tcp_hold_time_(XmppChannelConfig::kTcpHoldTime),
     228           7 :       gr_helper_disable_(false),
     229           7 :       dscp_value_(0),
     230           7 :       connection_queue_(TaskScheduler::GetInstance()->GetTaskId("bgp::Config"),
     231          14 :           0, boost::bind(&XmppServer::DequeueConnection, this, _1)) {
     232           7 : }
     233             : 
     234          40 : void XmppServer::CreateConfigUpdater(BgpConfigManager *config_manager) {
     235          40 :     xmpp_config_updater_.reset(new XmppConfigUpdater(this, config_manager));
     236          40 : }
     237             : 
     238          36 : uint16_t XmppServer::GetGracefulRestartTime() const {
     239             :     // Check if GR is disabled..
     240          36 :     if (!xmpp_config_updater_ || !xmpp_config_updater_->config().gr_enable())
     241           0 :         return 0;
     242          36 :     return xmpp_config_updater_->config().gr_time();
     243             : }
     244             : 
     245          16 : uint32_t XmppServer::GetLongLivedGracefulRestartTime() const {
     246             :     // Check if GR is disabled..
     247          16 :     if (!xmpp_config_updater_ || !xmpp_config_updater_->config().gr_enable())
     248           0 :         return 0;
     249          16 :     return xmpp_config_updater_->config().llgr_time();
     250             : }
     251             : 
     252         288 : uint32_t XmppServer::GetEndOfRibReceiveTime() const {
     253         288 :     if (xmpp_config_updater_)
     254         288 :         return xmpp_config_updater_->config().end_of_rib_timeout();
     255           0 :     return BgpGlobalSystemConfig::kEndOfRibTime;
     256             : }
     257             : 
     258        1196 : uint32_t XmppServer::GetEndOfRibSendTime() const {
     259        1196 :     if (xmpp_config_updater_)
     260        1196 :         xmpp_config_updater_->config().end_of_rib_timeout();
     261        1196 :     return BgpGlobalSystemConfig::kEndOfRibTime;
     262             : }
     263             : 
     264         449 : bool XmppServer::IsGRHelperModeEnabled() const {
     265             :     // Check if disabled in .conf file.
     266         449 :     if (gr_helper_disable_)
     267         128 :         return false;
     268             : 
     269             :     // Check from configuration.
     270         321 :     if (!xmpp_config_updater_)
     271          45 :         return false;
     272             : 
     273             :     // Check if GR is disabled..
     274         276 :     if (!xmpp_config_updater_->config().gr_enable())
     275          32 :         return false;
     276             : 
     277         244 :     return xmpp_config_updater_->config().gr_xmpp_helper();
     278             : }
     279             : 
     280         524 : bool XmppServer::IsPeerCloseGraceful() const {
     281             : 
     282             :     // If the server is deleted, do not do graceful restart
     283         524 :     if (deleter()->IsDeleted())
     284         144 :         return false;
     285             : 
     286             :     // Check if GR helper mode is disabled.
     287         380 :     if (!IsGRHelperModeEnabled())
     288         188 :         return false;
     289             : 
     290             :     // Enable GR if either gr-time or llgr-time is configured.
     291         192 :     return (xmpp_config_updater_->config().gr_time() ||
     292         192 :             xmpp_config_updater_->config().llgr_time());
     293             : }
     294             : 
     295         189 : XmppServer::~XmppServer() {
     296         120 :     STLDeleteElements(&connection_endpoint_map_);
     297         120 :     TcpServer::ClearSessions();
     298         189 : }
     299             : 
     300           0 : bool XmppServer::Initialize(short port) {
     301           0 :     log_uve_ = false;
     302           0 :     return TcpServer::Initialize(port);
     303             : }
     304             : 
     305         118 : bool XmppServer::Initialize(short port, bool logUVE) {
     306         118 :     log_uve_ = logUVE;
     307         118 :     return TcpServer::Initialize(port);
     308             : }
     309             : 
     310           0 : bool XmppServer::Initialize(short port, bool logUVE, const IpAddress& ip) {
     311           0 :     log_uve_ = logUVE;
     312           0 :     return TcpServer::Initialize(port, ip);
     313             : }
     314             : 
     315             : //
     316             : // Can be removed after Shutdown is renamed to ManagedDelete.
     317             : //
     318         118 : void XmppServer::SessionShutdown() {
     319         118 :     XmppConnectionManager::Shutdown();
     320         118 : }
     321             : 
     322             : //
     323             : // Return true if it's possible to delete the XmppServer.
     324             : //
     325             : // No need to check the connection WorkQueue since XmppServerConnections on it
     326             : // are dependents of the XmppServer.
     327             : //
     328         118 : bool XmppServer::MayDelete() const {
     329         118 :     return (GetSessionQueueSize() == 0);
     330             : }
     331             : 
     332             : //
     333             : // Trigger deletion of the XmppServer.
     334             : //
     335             : // A mutex is used to ensure that we do not create new XmppServerConnections
     336             : // after this point.  Note that this routine and AcceptSession may be called
     337             : // concurrently from 2 different threads in tests.
     338             : //
     339         118 : void XmppServer::Shutdown() {
     340         118 :     std::scoped_lock lock(deletion_mutex_);
     341         118 :     deleter_->Delete();
     342         118 : }
     343             : 
     344             : //
     345             : // Called when the XmppServer delete actor is being destroyed.
     346             : //
     347         118 : void XmppServer::Terminate() {
     348         118 :     ClearSessions();
     349         118 :     connection_queue_.Shutdown();
     350         118 : }
     351             : 
     352         302 : LifetimeActor *XmppServer::deleter() {
     353         302 :     return deleter_.get();
     354             : }
     355             : 
     356         524 : LifetimeActor *XmppServer::deleter() const {
     357         524 :     return deleter_.get();
     358             : }
     359             : 
     360         422 : LifetimeManager *XmppServer::lifetime_manager() {
     361         422 :     return lifetime_manager_.get();
     362             : }
     363             : 
     364          22 : TcpSession *XmppServer::CreateSession() {
     365             :     typedef boost::asio::detail::socket_option::boolean<
     366             :         SOL_SOCKET, SO_REUSEADDR> reuse_addr_t;
     367          22 :     TcpSession *session = TcpServer::CreateSession();
     368          22 :     Socket *socket = session->socket();
     369             : 
     370          22 :     boost::system::error_code err;
     371          22 :     socket->open(ip::tcp::v4(), err);
     372          22 :     if (err) {
     373           0 :         XMPP_WARNING(ServerOpenFail, err.message());
     374             :     }
     375             : 
     376          22 :     socket->set_option(reuse_addr_t(true), err);
     377          22 :     if (err) {
     378           0 :         XMPP_WARNING(SetSockOptFail, "", XMPP_PEER_DIR_OUT, err.message());
     379             :     }
     380             : 
     381          22 :     socket->bind(LocalEndpoint(), err);
     382          22 :     if (err) {
     383          22 :         XMPP_WARNING(ServerBindFailure, err.message());
     384             :     }
     385             : 
     386          22 :     XmppSession *xmpps = static_cast<XmppSession *>(session);
     387          22 :     err = xmpps->EnableTcpKeepalive(tcp_hold_time_);
     388          22 :     if (err) {
     389           0 :         XMPP_WARNING(ServerKeepAliveFailure, session->ToUVEKey(),
     390             :                      XMPP_PEER_DIR_OUT, err.message());
     391             :     }
     392             : 
     393          22 :     return session;
     394             : }
     395             : 
     396           0 : size_t XmppServer::ConnectionEventCount() const {
     397           0 :     return connection_event_map_.size();
     398             : }
     399             : 
     400           0 : size_t XmppServer::ConnectionMapSize() const {
     401           0 :     ReadLock lock(connection_map_mutex_);
     402           0 :     return connection_map_.size();
     403           0 : }
     404             : 
     405           0 : size_t XmppServer::ConnectionCount() const {
     406           0 :     return ConnectionMapSize() + deleted_connection_set_.size();
     407             : }
     408             : 
     409         278 : XmppServerConnection *XmppServer::FindConnection(Endpoint remote_endpoint) {
     410         278 :     ReadLock lock(connection_map_mutex_);
     411         278 :     ConnectionMap::iterator loc = connection_map_.find(remote_endpoint);
     412         278 :     if (loc != connection_map_.end()) {
     413           0 :         return loc->second;
     414             :     }
     415         278 :     return NULL;
     416         278 : }
     417             : 
     418         488 : XmppServerConnection *XmppServer::FindConnection(const string &address) {
     419         488 :     ReadLock lock(connection_map_mutex_);
     420         493 :     for (auto& value : connection_map_) {
     421          22 :         if (value.second->ToString() == address)
     422          17 :             return value.second;
     423             :     }
     424         471 :     return NULL;
     425         488 : }
     426             : 
     427           0 : bool XmppServer::ClearConnection(const string &hostname) {
     428           0 :     ReadLock lock(connection_map_mutex_);
     429           0 :     for (auto& value : connection_map_) {
     430           0 :         if (value.second->GetComputeHostName() == hostname) {
     431           0 :             value.second->Clear();
     432           0 :             return true;
     433             :         }
     434             :     }
     435           0 :     return false;
     436           0 : }
     437             : 
     438           0 : void XmppServer::UpdateAllConnections(uint8_t time_out) {
     439           0 :     ReadLock lock(connection_map_mutex_);
     440           0 :     for (auto& value : connection_map_) {
     441           0 :         value.second->UpdateKeepAliveTimer(time_out);
     442             :     }
     443           0 : }
     444             : 
     445         140 : void XmppServer::ClearAllConnections() {
     446         140 :     ReadLock lock(connection_map_mutex_);
     447         226 :     for (auto& value : connection_map_) {
     448          86 :         value.second->Clear();
     449             :     }
     450         140 : }
     451             : 
     452          49 : void XmppServer::RegisterConnectionEvent(xmps::PeerId id,
     453             :                                          ConnectionEventCb cb) {
     454          49 :     connection_event_map_.insert(make_pair(id, cb));
     455          49 : }
     456             : 
     457          40 : void XmppServer::UnRegisterConnectionEvent(xmps::PeerId id) {
     458          40 :     ConnectionEventCbMap::iterator it =  connection_event_map_.find(id);
     459          40 :     if (it != connection_event_map_.end())
     460          40 :         connection_event_map_.erase(it);
     461          40 : }
     462             : 
     463         576 : void XmppServer::NotifyConnectionEvent(XmppChannelMux *mux,
     464             :                                        xmps::PeerState state) {
     465         576 :     ConnectionEventCbMap::iterator iter = connection_event_map_.begin();
     466        1109 :     for (; iter != connection_event_map_.end(); ++iter) {
     467         533 :         ConnectionEventCb cb = iter->second;
     468         533 :         cb(mux, state);
     469         533 :     }
     470         576 : }
     471             : 
     472         300 : SslSession *XmppServer::AllocSession(SslSocket *socket) {
     473         300 :     SslSession *session = new XmppSession(this, socket);
     474         300 :     boost::system::error_code err;
     475         300 :     XmppSession *xmpp_session = static_cast<XmppSession *>(session);
     476         300 :     err = xmpp_session->EnableTcpKeepalive(tcp_hold_time_);
     477         300 :     if (err) {
     478          22 :         XMPP_WARNING(ServerKeepAliveFailure, xmpp_session->ToUVEKey(),
     479             :                      XMPP_PEER_DIR_OUT, err.message());
     480             :     }
     481         300 :     return session;
     482             : }
     483             : 
     484             : //
     485             : // Accept newly formed passive tcp session by creating necessary xmpp data
     486             : // structures. We do so to make sure that if there is any error reported
     487             : // over this tcp session, it can still be correctly handled, even though
     488             : // the allocated xmpp data structures are not fully processed yet.
     489             : //
     490         278 : bool XmppServer::AcceptSession(TcpSession *tcp_session) {
     491         278 :     std::scoped_lock lock(deletion_mutex_);
     492         278 :     if (deleter_->IsDeleted())
     493           0 :         return false;
     494             : 
     495         278 :     XmppSession *session = dynamic_cast<XmppSession *>(tcp_session);
     496         278 :     XmppServerConnection *connection = CreateConnection(session);
     497             : 
     498         278 :     if (xmpp_config_updater_) {
     499         530 :         connection->state_machine()->set_hold_time(
     500         265 :                 xmpp_config_updater_->xmpp_hold_time());
     501             :     }
     502             :     // Register event handler.
     503         278 :     tcp_session->set_observer(boost::bind(&XmppStateMachine::OnSessionEvent,
     504             :                                           connection->state_machine(), _1, _2));
     505             : 
     506             :     // set async_read_ready as false
     507         278 :     session->set_read_on_connect(false);
     508         278 :     connection->set_session(session);
     509         278 :     connection->state_machine()->set_session(session);
     510         278 :     connection->set_on_work_queue();
     511         278 :     connection_queue_.Enqueue(connection);
     512         278 :     return true;
     513         278 : }
     514             : 
     515             : //
     516             : // Remove the given XmppServerConnection from the ConnectionMap.
     517             : //
     518         302 : void XmppServer::RemoveConnection(XmppServerConnection *connection) {
     519         302 :     CHECK_CONCURRENCY("bgp::Config");
     520             : 
     521         302 :     assert(connection->IsDeleted());
     522         302 :     Endpoint endpoint = connection->endpoint();
     523         302 :     WriteLock lock(connection_map_mutex_);
     524         302 :     ConnectionMap::iterator loc = connection_map_.find(endpoint);
     525         302 :     assert(loc != connection_map_.end() && loc->second == connection);
     526         302 :     connection_map_.erase(loc);
     527         302 : }
     528             : 
     529          20 : void XmppServer::SwapXmppConnectionMapEntries(
     530             :         XmppConnection *connection1, XmppConnection *connection2) {
     531          20 :     WriteLock lock(connection_map_mutex_);
     532             :     ConnectionMap::iterator loc1 =
     533          20 :         connection_map_.find(connection1->endpoint());
     534          20 :     assert(loc1 != connection_map_.end());
     535             :     ConnectionMap::iterator loc2 =
     536          20 :         connection_map_.find(connection2->endpoint());
     537          20 :     assert(loc2 != connection_map_.end());
     538          20 :     swap(loc1->second, loc2->second);
     539          20 :     swap(loc1->second->endpoint(), loc2->second->endpoint());
     540          20 : }
     541             : 
     542             : //
     543             : // Insert the given XmppServerConnection into the ConnectionMap.
     544             : //
     545         302 : void XmppServer::InsertConnection(XmppServerConnection *connection) {
     546         302 :     CHECK_CONCURRENCY("bgp::Config");
     547             : 
     548         302 :     assert(!connection->IsDeleted());
     549         302 :     connection->Initialize();
     550         302 :     Endpoint endpoint = connection->endpoint();
     551         302 :     ConnectionMap::iterator loc;
     552             :     bool result;
     553         302 :     WriteLock lock(connection_map_mutex_);
     554         302 :     tie(loc, result) = connection_map_.insert(make_pair(endpoint, connection));
     555         302 :     assert(result);
     556         302 :     max_connections_ = max(max_connections_, connection_map_.size());
     557         302 : }
     558             : 
     559             : //
     560             : // Create XmppConnnection and its associated data structures. This API is
     561             : // only used to allocate data structures and initialize necessary fields.
     562             : // The data structures are not populated to any maps in the XmppServer at
     563             : // this point.  However, the newly created XmppServerConnection does add
     564             : // itself as a dependent of the XmppServer via LifetimeManager linkage.
     565             : //
     566         278 : XmppServerConnection *XmppServer::CreateConnection(XmppSession *session) {
     567             :     XmppServerConnection  *connection;
     568             : 
     569         278 :     XmppChannelConfig cfg(false);
     570         278 :     cfg.endpoint = session->remote_endpoint();
     571         278 :     cfg.local_endpoint = session->local_endpoint();
     572         278 :     cfg.FromAddr = server_addr_;
     573         278 :     cfg.logUVE = log_uve_;
     574         278 :     cfg.auth_enabled = auth_enabled_;
     575         278 :     cfg.dscp_value = dscp_value_;
     576             : 
     577         278 :     XMPP_DEBUG(XmppCreateConnection, session->ToUVEKey(), XMPP_PEER_DIR_OUT,
     578             :                session->ToString());
     579             :     connection = XmppStaticObjectFactory::Create<XmppServerConnection>
     580         278 :         (this, static_cast<const XmppChannelConfig*>(&cfg));
     581             : 
     582         278 :     return connection;
     583         278 : }
     584             : 
     585           0 : void XmppServer::SetDscpValue(uint8_t value) {
     586           0 :     ReadLock lock(connection_map_mutex_);
     587           0 :     dscp_value_ = value;
     588           0 :     for (auto& value : connection_map_) {
     589           0 :         XmppServerConnection *connection = value.second;
     590           0 :         connection->SetDscpValue(dscp_value_);
     591             :     }
     592           0 : }
     593             : //
     594             : // Handler for XmppServerConnections that are dequeued from the WorkQueue.
     595             : //
     596             : // Since the XmppServerConnections on the WorkQueue are dependents of the
     597             : // XmppServer, we are guaranteed that the XmppServer won't get destroyed
     598             : // before the connection WorkQueue is drained.
     599             : //
     600         278 : bool XmppServer::DequeueConnection(XmppServerConnection *connection) {
     601         278 :     CHECK_CONCURRENCY("bgp::Config");
     602         278 :     connection->clear_on_work_queue();
     603             : 
     604             :     // This happens if the XmppServer got deleted while the XmppConnnection
     605             :     // was on the WorkQueue.
     606         278 :     if (connection->IsDeleted()) {
     607           0 :         connection->RetryDelete();
     608           0 :         return true;
     609             :     }
     610             : 
     611         278 :     XmppSession *session = connection->session();
     612         278 :     connection->clear_session();
     613         278 :     Endpoint remote_endpoint = session->remote_endpoint();
     614         278 :     XmppServerConnection *old_connection = FindConnection(remote_endpoint);
     615             : 
     616             :     // Close as duplicate if we have a connection from the same Endpoint.
     617             :     // Otherwise go ahead and insert into the ConnectionMap. We may find
     618             :     // it has a conflicting Endpoint name and decide to terminate it when
     619             :     // we process the Open message.
     620         278 :     if (old_connection) {
     621           0 :         XMPP_DEBUG(XmppCreateConnection, session->ToUVEKey(), XMPP_PEER_DIR_IN,
     622             :                    "Close duplicate connection " + session->ToString());
     623             :         // Remove reference to the session from StateMachine directly. We don't
     624             :         // go through the entire normal cleanup pipeline for cleaning these
     625             :         // duplicate connections.
     626           0 :         assert(connection->state_machine()->session());
     627           0 :         assert(connection->state_machine()->session() == session);
     628           0 :         connection->state_machine()->RemoveSession();
     629           0 :         DeleteSession(session);
     630           0 :         connection->set_duplicate();
     631           0 :         connection->ManagedDelete();
     632           0 :         InsertDeletedConnection(connection);
     633             :     } else {
     634         278 :         InsertConnection(connection);
     635         278 :         connection->AcceptSession(session);
     636             :     }
     637             : 
     638         278 :     return true;
     639             : }
     640             : 
     641           0 : size_t XmppServer::GetConnectionQueueSize() const {
     642           0 :     return connection_queue_.Length();
     643             : }
     644             : 
     645           0 : void XmppServer::SetConnectionQueueDisable(bool disabled) {
     646           0 :     connection_queue_.set_disable(disabled);
     647           0 : }
     648             : 
     649             : //
     650             : // Connection is marked deleted, add it to the ConnectionSet.
     651             : //
     652         302 : void XmppServer::InsertDeletedConnection(XmppServerConnection *connection) {
     653         302 :     CHECK_CONCURRENCY("bgp::Config");
     654             : 
     655         302 :     assert(connection->IsDeleted());
     656         302 :     ConnectionSet::iterator it;
     657             :     bool result;
     658         302 :     tie(it, result) = deleted_connection_set_.insert(connection);
     659         302 :     assert(result);
     660         302 : }
     661             : 
     662             : //
     663             : // Connection is being destroyed, remove it from the ConnectionSet.
     664             : //
     665         302 : void XmppServer::RemoveDeletedConnection(XmppServerConnection *connection) {
     666         302 :     CHECK_CONCURRENCY("bgp::Config");
     667             : 
     668         302 :     assert(connection->IsDeleted());
     669         302 :     ConnectionSet::iterator it = deleted_connection_set_.find(connection);
     670         302 :     assert(it != deleted_connection_set_.end());
     671         302 :     deleted_connection_set_.erase(it);
     672         302 :     ReleaseConnectionEndpoint(connection);
     673         302 : }
     674             : 
     675         573 : XmppConnectionEndpoint *XmppServer::FindConnectionEndpoint(
     676             :     const string &endpoint_name) {
     677         573 :     std::scoped_lock lock(endpoint_map_mutex_);
     678             :     ConnectionEndpointMap::const_iterator loc =
     679         573 :         connection_endpoint_map_.find(endpoint_name);
     680        1146 :     return (loc != connection_endpoint_map_.end() ? loc->second : NULL);
     681         573 : }
     682             : 
     683         275 : XmppConnectionEndpoint *XmppServer::LocateConnectionEndpoint(
     684             :         XmppServerConnection *connection, bool &created) {
     685         275 :     created = false;
     686         275 :     if (!connection)
     687           0 :         return NULL;
     688             : 
     689         275 :     std::scoped_lock lock(endpoint_map_mutex_);
     690             : 
     691             :     ConnectionEndpointMap::const_iterator loc =
     692         275 :         connection_endpoint_map_.find(connection->ToString());
     693             :     XmppConnectionEndpoint *conn_endpoint;
     694             : 
     695         275 :     if (loc != connection_endpoint_map_.end()) {
     696          92 :         conn_endpoint = loc->second;
     697          92 :         if (!conn_endpoint->connection()) {
     698          72 :             created = true;
     699          72 :             conn_endpoint->set_connection(connection);
     700          72 :             connection->set_conn_endpoint(conn_endpoint);
     701             :         }
     702          92 :         return conn_endpoint;
     703             :     }
     704             : 
     705         183 :     created = true;
     706         183 :     conn_endpoint = new XmppConnectionEndpoint(connection->ToString());
     707             :     bool result;
     708         183 :     tie(loc, result) = connection_endpoint_map_.insert(
     709         366 :             make_pair(connection->ToString(), conn_endpoint));
     710         183 :     assert(result);
     711         183 :     conn_endpoint->set_connection(connection);
     712         183 :     connection->set_conn_endpoint(conn_endpoint);
     713         183 :     return conn_endpoint;
     714         275 : }
     715             : 
     716             : //
     717             : // Remove association of the given XmppConnectionEndpoint with XmppConnection.
     718             : // This method is provided just to make things symmetrical - the caller could
     719             : // simply have called XmppConnectionEndpoint::reset_connection directly.
     720             : //
     721         604 : void XmppServer::ReleaseConnectionEndpoint(XmppServerConnection *connection) {
     722         604 :     std::scoped_lock lock(endpoint_map_mutex_);
     723             : 
     724         604 :     if (!connection->conn_endpoint())
     725         349 :         return;
     726         255 :     assert(connection->conn_endpoint()->connection() == connection);
     727         255 :     connection->conn_endpoint()->reset_connection();
     728         255 :     connection->set_conn_endpoint(NULL);
     729         604 : }
     730             : 
     731           0 : void XmppServer::FillShowConnections(
     732             :     vector<ShowXmppConnection> *show_connection_list) const {
     733           0 :     ReadLock lock(connection_map_mutex_);
     734           0 :     for (const auto& value : connection_map_) {
     735           0 :         const XmppServerConnection *connection = value.second;
     736           0 :         ShowXmppConnection show_connection;
     737           0 :         connection->FillShowInfo(&show_connection);
     738           0 :         show_connection_list->push_back(show_connection);
     739           0 :     }
     740           0 :     for (const auto& connection : deleted_connection_set_) {
     741           0 :         ShowXmppConnection show_connection;
     742           0 :         connection->FillShowInfo(&show_connection);
     743           0 :         show_connection_list->push_back(show_connection);
     744           0 :     }
     745           0 : }
     746             : 
     747           0 : void XmppServer::FillShowServer(ShowXmppServerResp *resp) const {
     748           0 :     SocketIOStats peer_socket_stats;
     749           0 :     GetRxSocketStats(&peer_socket_stats);
     750           0 :     resp->set_rx_socket_stats(peer_socket_stats);
     751           0 :     GetTxSocketStats(&peer_socket_stats);
     752           0 :     resp->set_tx_socket_stats(peer_socket_stats);
     753           0 :     resp->set_current_connections(ConnectionMapSize());
     754           0 :     resp->set_max_connections(max_connections_);
     755           0 : }
     756             : 
     757             : class ShowXmppConnectionHandler {
     758             : public:
     759           0 :     static bool CallbackS1(const Sandesh *sr,
     760             :             const RequestPipeline::PipeSpec ps, int stage, int instNum,
     761             :             RequestPipeline::InstData *data) {
     762             :         const ShowXmppConnectionReq *req =
     763           0 :             static_cast<const ShowXmppConnectionReq *>(ps.snhRequest_.get());
     764             :         XmppSandeshContext *xsc =
     765           0 :             dynamic_cast<XmppSandeshContext *>(req->module_context("XMPP"));
     766             : 
     767           0 :         ShowXmppConnectionResp *resp = new ShowXmppConnectionResp;
     768           0 :         vector<ShowXmppConnection> connections;
     769           0 :         if (xsc)
     770           0 :             xsc->xmpp_server->FillShowConnections(&connections);
     771           0 :         resp->set_connections(connections);
     772           0 :         resp->set_context(req->context());
     773           0 :         resp->Response();
     774           0 :         return true;
     775           0 :     }
     776             : };
     777             : 
     778           0 : void ShowXmppConnectionReq::HandleRequest() const {
     779           0 :     RequestPipeline::PipeSpec ps(this);
     780             : 
     781             :     // Request pipeline has single stage to collect connection info and
     782             :     // respond to the request.
     783           0 :     RequestPipeline::StageSpec s1;
     784           0 :     TaskScheduler *scheduler = TaskScheduler::GetInstance();
     785           0 :     s1.taskId_ = scheduler->GetTaskId("bgp::ShowCommand");
     786           0 :     s1.cbFn_ = ShowXmppConnectionHandler::CallbackS1;
     787           0 :     s1.instances_.push_back(0);
     788           0 :     ps.stages_.push_back(s1);
     789           0 :     RequestPipeline rp(ps);
     790           0 : }
     791             : 
     792             : class ClearXmppConnectionHandler {
     793             : public:
     794           0 :     static bool CallbackS1(const Sandesh *sr,
     795             :         const RequestPipeline::PipeSpec ps, int stage, int instNum,
     796             :         RequestPipeline::InstData *data) {
     797             : 
     798             :         const ClearXmppConnectionReq *req =
     799           0 :             static_cast<const ClearXmppConnectionReq *>(ps.snhRequest_.get());
     800             :         XmppSandeshContext *xsc =
     801           0 :             dynamic_cast<XmppSandeshContext *>(req->module_context("XMPP"));
     802             : 
     803           0 :         ClearXmppConnectionResp *resp = new ClearXmppConnectionResp;
     804           0 :         if (!xsc || !xsc->test_mode) {
     805           0 :             resp->set_success(false);
     806           0 :         } else if (req->get_hostname_or_all() != "all") {
     807           0 :             if (xsc->xmpp_server->ClearConnection(req->get_hostname_or_all())) {
     808           0 :                 resp->set_success(true);
     809             :             } else {
     810           0 :                 resp->set_success(false);
     811             :             }
     812             :         } else {
     813           0 :             if (xsc->xmpp_server->ConnectionCount()) {
     814           0 :                 xsc->xmpp_server->ClearAllConnections();
     815           0 :                 resp->set_success(true);
     816             :             } else {
     817           0 :                 resp->set_success(false);
     818             :             }
     819             :         }
     820             : 
     821           0 :         resp->set_context(req->context());
     822           0 :         resp->Response();
     823           0 :         return true;
     824             :     }
     825             : };
     826             : 
     827           0 : void ClearXmppConnectionReq::HandleRequest() const {
     828             : 
     829             :     // config task is used to create and delete connection objects.
     830             :     // hence use the same task to find the connection
     831           0 :     RequestPipeline::StageSpec s1;
     832           0 :     s1.taskId_ = TaskScheduler::GetInstance()->GetTaskId("bgp::Config");
     833           0 :     s1.instances_.push_back(0);
     834           0 :     s1.cbFn_ = ClearXmppConnectionHandler::CallbackS1;
     835             : 
     836           0 :     RequestPipeline::PipeSpec ps(this);
     837           0 :     ps.stages_.push_back(s1);
     838           0 :     RequestPipeline rp(ps);
     839           0 : }
     840             : 
     841             : class ShowXmppServerHandler {
     842             : public:
     843           0 :     static bool CallbackS1(const Sandesh *sr,
     844             :             const RequestPipeline::PipeSpec ps, int stage, int instNum,
     845             :             RequestPipeline::InstData *data) {
     846             :         const ShowXmppServerReq *req =
     847           0 :             static_cast<const ShowXmppServerReq *>(ps.snhRequest_.get());
     848             :         XmppSandeshContext *xsc =
     849           0 :             dynamic_cast<XmppSandeshContext *>(req->module_context("XMPP"));
     850             : 
     851           0 :         ShowXmppServerResp *resp = new ShowXmppServerResp;
     852           0 :         if (xsc)
     853           0 :             xsc->xmpp_server->FillShowServer(resp);
     854           0 :         resp->set_context(req->context());
     855           0 :         resp->Response();
     856           0 :         return true;
     857             :     }
     858             : };
     859             : 
     860           0 : void ShowXmppServerReq::HandleRequest() const {
     861           0 :     RequestPipeline::PipeSpec ps(this);
     862           0 :     RequestPipeline::StageSpec s1;
     863           0 :     TaskScheduler *scheduler = TaskScheduler::GetInstance();
     864           0 :     s1.taskId_ = scheduler->GetTaskId("bgp::ShowCommand");
     865           0 :     s1.cbFn_ = ShowXmppServerHandler::CallbackS1;
     866           0 :     s1.instances_.push_back(0);
     867           0 :     ps.stages_.push_back(s1);
     868           0 :     RequestPipeline rp(ps);
     869           0 : }

Generated by: LCOV version 1.14