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 80 : DeleteActor(XmppServer *server)
33 80 : : LifetimeActor(server->lifetime_manager()), server_(server) {
34 80 : }
35 78 : virtual bool MayDelete() const {
36 78 : CHECK_CONCURRENCY("bgp::Config");
37 78 : return server_->MayDelete();
38 : }
39 78 : virtual void Shutdown() {
40 78 : CHECK_CONCURRENCY("bgp::Config");
41 78 : server_->SessionShutdown();
42 78 : }
43 78 : virtual void Destroy() {
44 78 : CHECK_CONCURRENCY("bgp::Config");
45 78 : server_->Terminate();
46 78 : }
47 :
48 : private:
49 : XmppServer *server_;
50 : };
51 :
52 :
53 34 : XmppServer::XmppServer(EventManager *evm, const string &server_addr,
54 34 : const XmppChannelConfig *config)
55 : : XmppConnectionManager(
56 34 : evm, ssl::context::sslv23_server, config->auth_enabled, true),
57 34 : max_connections_(0),
58 34 : lifetime_manager_(XmppStaticObjectFactory::Create<XmppLifetimeManager>(
59 : TaskScheduler::GetInstance()->GetTaskId("bgp::Config"))),
60 34 : deleter_(new DeleteActor(this)),
61 34 : server_addr_(server_addr),
62 34 : log_uve_(false),
63 34 : auth_enabled_(config->auth_enabled),
64 34 : tcp_hold_time_(config->tcp_hold_time),
65 34 : gr_helper_disable_(config->gr_helper_disable),
66 34 : dscp_value_(0),
67 34 : connection_queue_(TaskScheduler::GetInstance()->GetTaskId("bgp::Config"),
68 102 : 0, boost::bind(&XmppServer::DequeueConnection, this, _1)) {
69 :
70 34 : if (config->auth_enabled) {
71 :
72 : // Get SSL context from base class and update
73 34 : boost::asio::ssl::context *ctx = context();
74 34 : boost::system::error_code ec;
75 :
76 : // set mode
77 34 : ctx->set_options(ssl::context::default_workarounds |
78 : ssl::context::no_sslv3 | ssl::context::no_sslv2, ec);
79 34 : 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 34 : std::string ca_cert_filename = config->path_to_ca_cert;
87 34 : 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 34 : ctx->use_certificate_file(config->path_to_server_cert,
108 : boost::asio::ssl::context::pem, ec);
109 34 : 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 34 : ctx->use_private_key_file(config->path_to_server_priv_key,
118 : boost::asio::ssl::context::pem, ec);
119 34 : 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 34 : }
126 34 : }
127 :
128 : class XmppConfigUpdater {
129 : public:
130 0 : explicit XmppConfigUpdater(XmppServer *server,
131 0 : BgpConfigManager *config_manager) :
132 0 : server_(server) {
133 0 : BgpConfigManager::Observers obs;
134 : obs.system= boost::bind(&XmppConfigUpdater::ProcessGlobalSystemConfig,
135 0 : this, _1, _2);
136 : obs.protocol = boost::bind(&XmppConfigUpdater::ProcessProtocolConfig,
137 0 : this, _1, _2);
138 0 : config_manager->RegisterObservers(obs);
139 0 : }
140 :
141 0 : const BgpGlobalSystemConfig &config() const { return config_; }
142 :
143 0 : void ProcessProtocolConfig(const BgpProtocolConfig *protocol_config,
144 : BgpConfigManager::EventType event) {
145 :
146 0 : if (server_->subcluster_name() != protocol_config->subcluster_name()) {
147 0 : server_->set_subcluster_name(protocol_config->subcluster_name());
148 0 : server_->ClearAllConnections();
149 : }
150 0 : }
151 :
152 0 : void ProcessGlobalSystemConfig(const BgpGlobalSystemConfig *system,
153 : BgpConfigManager::EventType event) {
154 : // Clear peers only if GR is or was enabled.
155 0 : bool clear_peers = config_.gr_enable() || system->gr_enable();
156 0 : bool update_peers = false;
157 0 : config_.set_gr_enable(system->gr_enable());
158 0 : config_.set_gr_time(system->gr_time());
159 0 : config_.set_llgr_time(system->llgr_time());
160 0 : config_.set_end_of_rib_timeout(system->end_of_rib_timeout());
161 0 : config_.set_gr_xmpp_helper(system->gr_xmpp_helper());
162 :
163 : // Process any change in xmpp-hold-time
164 0 : if (config_.fc_enabled() != system->fc_enabled() ||
165 0 : 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 0 : 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 0 : if (clear_peers)
181 0 : server_->ClearAllConnections();
182 0 : else if (update_peers)
183 0 : server_->UpdateAllConnections(config_.xmpp_hold_time());
184 0 : }
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 0 : 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 0 : void XmppServer::CreateConfigUpdater(BgpConfigManager *config_manager) {
235 0 : xmpp_config_updater_.reset(new XmppConfigUpdater(this, config_manager));
236 0 : }
237 :
238 0 : uint16_t XmppServer::GetGracefulRestartTime() const {
239 : // Check if GR is disabled..
240 0 : if (!xmpp_config_updater_ || !xmpp_config_updater_->config().gr_enable())
241 0 : return 0;
242 0 : return xmpp_config_updater_->config().gr_time();
243 : }
244 :
245 0 : uint32_t XmppServer::GetLongLivedGracefulRestartTime() const {
246 : // Check if GR is disabled..
247 0 : if (!xmpp_config_updater_ || !xmpp_config_updater_->config().gr_enable())
248 0 : return 0;
249 0 : return xmpp_config_updater_->config().llgr_time();
250 : }
251 :
252 0 : uint32_t XmppServer::GetEndOfRibReceiveTime() const {
253 0 : if (xmpp_config_updater_)
254 0 : return xmpp_config_updater_->config().end_of_rib_timeout();
255 0 : return BgpGlobalSystemConfig::kEndOfRibTime;
256 : }
257 :
258 0 : uint32_t XmppServer::GetEndOfRibSendTime() const {
259 0 : if (xmpp_config_updater_)
260 0 : xmpp_config_updater_->config().end_of_rib_timeout();
261 0 : return BgpGlobalSystemConfig::kEndOfRibTime;
262 : }
263 :
264 45 : bool XmppServer::IsGRHelperModeEnabled() const {
265 : // Check if disabled in .conf file.
266 45 : if (gr_helper_disable_)
267 0 : return false;
268 :
269 : // Check from configuration.
270 45 : if (!xmpp_config_updater_)
271 45 : return false;
272 :
273 : // Check if GR is disabled..
274 0 : if (!xmpp_config_updater_->config().gr_enable())
275 0 : return false;
276 :
277 0 : return xmpp_config_updater_->config().gr_xmpp_helper();
278 : }
279 :
280 28 : bool XmppServer::IsPeerCloseGraceful() const {
281 :
282 : // If the server is deleted, do not do graceful restart
283 28 : if (deleter()->IsDeleted())
284 0 : return false;
285 :
286 : // Check if GR helper mode is disabled.
287 28 : if (!IsGRHelperModeEnabled())
288 28 : return false;
289 :
290 : // Enable GR if either gr-time or llgr-time is configured.
291 0 : return (xmpp_config_updater_->config().gr_time() ||
292 0 : xmpp_config_updater_->config().llgr_time());
293 : }
294 :
295 149 : XmppServer::~XmppServer() {
296 80 : STLDeleteElements(&connection_endpoint_map_);
297 80 : TcpServer::ClearSessions();
298 149 : }
299 :
300 0 : bool XmppServer::Initialize(short port) {
301 0 : log_uve_ = false;
302 0 : return TcpServer::Initialize(port);
303 : }
304 :
305 78 : bool XmppServer::Initialize(short port, bool logUVE) {
306 78 : log_uve_ = logUVE;
307 78 : 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 78 : void XmppServer::SessionShutdown() {
319 78 : XmppConnectionManager::Shutdown();
320 78 : }
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 78 : bool XmppServer::MayDelete() const {
329 78 : 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 78 : void XmppServer::Shutdown() {
340 78 : std::scoped_lock lock(deletion_mutex_);
341 78 : deleter_->Delete();
342 78 : }
343 :
344 : //
345 : // Called when the XmppServer delete actor is being destroyed.
346 : //
347 78 : void XmppServer::Terminate() {
348 78 : ClearSessions();
349 78 : connection_queue_.Shutdown();
350 78 : }
351 :
352 37 : LifetimeActor *XmppServer::deleter() {
353 37 : return deleter_.get();
354 : }
355 :
356 28 : LifetimeActor *XmppServer::deleter() const {
357 28 : return deleter_.get();
358 : }
359 :
360 117 : LifetimeManager *XmppServer::lifetime_manager() {
361 117 : 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 13 : XmppServerConnection *XmppServer::FindConnection(Endpoint remote_endpoint) {
410 13 : ReadLock lock(connection_map_mutex_);
411 13 : ConnectionMap::iterator loc = connection_map_.find(remote_endpoint);
412 13 : if (loc != connection_map_.end()) {
413 0 : return loc->second;
414 : }
415 13 : return NULL;
416 13 : }
417 :
418 495 : XmppServerConnection *XmppServer::FindConnection(const string &address) {
419 495 : ReadLock lock(connection_map_mutex_);
420 505 : for (auto& value : connection_map_) {
421 27 : if (value.second->ToString() == address)
422 17 : return value.second;
423 : }
424 478 : return NULL;
425 495 : }
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 0 : void XmppServer::ClearAllConnections() {
446 0 : ReadLock lock(connection_map_mutex_);
447 0 : for (auto& value : connection_map_) {
448 0 : value.second->Clear();
449 : }
450 0 : }
451 :
452 9 : void XmppServer::RegisterConnectionEvent(xmps::PeerId id,
453 : ConnectionEventCb cb) {
454 9 : connection_event_map_.insert(make_pair(id, cb));
455 9 : }
456 :
457 0 : void XmppServer::UnRegisterConnectionEvent(xmps::PeerId id) {
458 0 : ConnectionEventCbMap::iterator it = connection_event_map_.find(id);
459 0 : if (it != connection_event_map_.end())
460 0 : connection_event_map_.erase(it);
461 0 : }
462 :
463 51 : void XmppServer::NotifyConnectionEvent(XmppChannelMux *mux,
464 : xmps::PeerState state) {
465 51 : ConnectionEventCbMap::iterator iter = connection_event_map_.begin();
466 59 : for (; iter != connection_event_map_.end(); ++iter) {
467 8 : ConnectionEventCb cb = iter->second;
468 8 : cb(mux, state);
469 8 : }
470 51 : }
471 :
472 35 : SslSession *XmppServer::AllocSession(SslSocket *socket) {
473 35 : SslSession *session = new XmppSession(this, socket);
474 35 : boost::system::error_code err;
475 35 : XmppSession *xmpp_session = static_cast<XmppSession *>(session);
476 35 : err = xmpp_session->EnableTcpKeepalive(tcp_hold_time_);
477 35 : if (err) {
478 22 : XMPP_WARNING(ServerKeepAliveFailure, xmpp_session->ToUVEKey(),
479 : XMPP_PEER_DIR_OUT, err.message());
480 : }
481 35 : 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 13 : bool XmppServer::AcceptSession(TcpSession *tcp_session) {
491 13 : std::scoped_lock lock(deletion_mutex_);
492 13 : if (deleter_->IsDeleted())
493 0 : return false;
494 :
495 13 : XmppSession *session = dynamic_cast<XmppSession *>(tcp_session);
496 13 : XmppServerConnection *connection = CreateConnection(session);
497 :
498 13 : if (xmpp_config_updater_) {
499 0 : connection->state_machine()->set_hold_time(
500 0 : xmpp_config_updater_->xmpp_hold_time());
501 : }
502 : // Register event handler.
503 13 : tcp_session->set_observer(boost::bind(&XmppStateMachine::OnSessionEvent,
504 : connection->state_machine(), _1, _2));
505 :
506 : // set async_read_ready as false
507 13 : session->set_read_on_connect(false);
508 13 : connection->set_session(session);
509 13 : connection->state_machine()->set_session(session);
510 13 : connection->set_on_work_queue();
511 13 : connection_queue_.Enqueue(connection);
512 13 : return true;
513 13 : }
514 :
515 : //
516 : // Remove the given XmppServerConnection from the ConnectionMap.
517 : //
518 37 : void XmppServer::RemoveConnection(XmppServerConnection *connection) {
519 37 : CHECK_CONCURRENCY("bgp::Config");
520 :
521 37 : assert(connection->IsDeleted());
522 37 : Endpoint endpoint = connection->endpoint();
523 37 : WriteLock lock(connection_map_mutex_);
524 37 : ConnectionMap::iterator loc = connection_map_.find(endpoint);
525 37 : assert(loc != connection_map_.end() && loc->second == connection);
526 37 : connection_map_.erase(loc);
527 37 : }
528 :
529 0 : void XmppServer::SwapXmppConnectionMapEntries(
530 : XmppConnection *connection1, XmppConnection *connection2) {
531 0 : WriteLock lock(connection_map_mutex_);
532 : ConnectionMap::iterator loc1 =
533 0 : connection_map_.find(connection1->endpoint());
534 0 : assert(loc1 != connection_map_.end());
535 : ConnectionMap::iterator loc2 =
536 0 : connection_map_.find(connection2->endpoint());
537 0 : assert(loc2 != connection_map_.end());
538 0 : swap(loc1->second, loc2->second);
539 0 : swap(loc1->second->endpoint(), loc2->second->endpoint());
540 0 : }
541 :
542 : //
543 : // Insert the given XmppServerConnection into the ConnectionMap.
544 : //
545 37 : void XmppServer::InsertConnection(XmppServerConnection *connection) {
546 37 : CHECK_CONCURRENCY("bgp::Config");
547 :
548 37 : assert(!connection->IsDeleted());
549 37 : connection->Initialize();
550 37 : Endpoint endpoint = connection->endpoint();
551 37 : ConnectionMap::iterator loc;
552 : bool result;
553 37 : WriteLock lock(connection_map_mutex_);
554 37 : tie(loc, result) = connection_map_.insert(make_pair(endpoint, connection));
555 37 : assert(result);
556 37 : max_connections_ = max(max_connections_, connection_map_.size());
557 37 : }
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 13 : XmppServerConnection *XmppServer::CreateConnection(XmppSession *session) {
567 : XmppServerConnection *connection;
568 :
569 13 : XmppChannelConfig cfg(false);
570 13 : cfg.endpoint = session->remote_endpoint();
571 13 : cfg.local_endpoint = session->local_endpoint();
572 13 : cfg.FromAddr = server_addr_;
573 13 : cfg.logUVE = log_uve_;
574 13 : cfg.auth_enabled = auth_enabled_;
575 13 : cfg.dscp_value = dscp_value_;
576 :
577 13 : XMPP_DEBUG(XmppCreateConnection, session->ToUVEKey(), XMPP_PEER_DIR_OUT,
578 : session->ToString());
579 : connection = XmppStaticObjectFactory::Create<XmppServerConnection>
580 13 : (this, static_cast<const XmppChannelConfig*>(&cfg));
581 :
582 13 : return connection;
583 13 : }
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 13 : bool XmppServer::DequeueConnection(XmppServerConnection *connection) {
601 13 : CHECK_CONCURRENCY("bgp::Config");
602 13 : connection->clear_on_work_queue();
603 :
604 : // This happens if the XmppServer got deleted while the XmppConnnection
605 : // was on the WorkQueue.
606 13 : if (connection->IsDeleted()) {
607 0 : connection->RetryDelete();
608 0 : return true;
609 : }
610 :
611 13 : XmppSession *session = connection->session();
612 13 : connection->clear_session();
613 13 : Endpoint remote_endpoint = session->remote_endpoint();
614 13 : 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 13 : 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 13 : InsertConnection(connection);
635 13 : connection->AcceptSession(session);
636 : }
637 :
638 13 : 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 37 : void XmppServer::InsertDeletedConnection(XmppServerConnection *connection) {
653 37 : CHECK_CONCURRENCY("bgp::Config");
654 :
655 37 : assert(connection->IsDeleted());
656 37 : ConnectionSet::iterator it;
657 : bool result;
658 37 : tie(it, result) = deleted_connection_set_.insert(connection);
659 37 : assert(result);
660 37 : }
661 :
662 : //
663 : // Connection is being destroyed, remove it from the ConnectionSet.
664 : //
665 37 : void XmppServer::RemoveDeletedConnection(XmppServerConnection *connection) {
666 37 : CHECK_CONCURRENCY("bgp::Config");
667 :
668 37 : assert(connection->IsDeleted());
669 37 : ConnectionSet::iterator it = deleted_connection_set_.find(connection);
670 37 : assert(it != deleted_connection_set_.end());
671 37 : deleted_connection_set_.erase(it);
672 37 : ReleaseConnectionEndpoint(connection);
673 37 : }
674 :
675 38 : XmppConnectionEndpoint *XmppServer::FindConnectionEndpoint(
676 : const string &endpoint_name) {
677 38 : std::scoped_lock lock(endpoint_map_mutex_);
678 : ConnectionEndpointMap::const_iterator loc =
679 38 : connection_endpoint_map_.find(endpoint_name);
680 76 : return (loc != connection_endpoint_map_.end() ? loc->second : NULL);
681 38 : }
682 :
683 23 : XmppConnectionEndpoint *XmppServer::LocateConnectionEndpoint(
684 : XmppServerConnection *connection, bool &created) {
685 23 : created = false;
686 23 : if (!connection)
687 0 : return NULL;
688 :
689 23 : std::scoped_lock lock(endpoint_map_mutex_);
690 :
691 : ConnectionEndpointMap::const_iterator loc =
692 23 : connection_endpoint_map_.find(connection->ToString());
693 : XmppConnectionEndpoint *conn_endpoint;
694 :
695 23 : if (loc != connection_endpoint_map_.end()) {
696 0 : conn_endpoint = loc->second;
697 0 : if (!conn_endpoint->connection()) {
698 0 : created = true;
699 0 : conn_endpoint->set_connection(connection);
700 0 : connection->set_conn_endpoint(conn_endpoint);
701 : }
702 0 : return conn_endpoint;
703 : }
704 :
705 23 : created = true;
706 23 : conn_endpoint = new XmppConnectionEndpoint(connection->ToString());
707 : bool result;
708 23 : tie(loc, result) = connection_endpoint_map_.insert(
709 46 : make_pair(connection->ToString(), conn_endpoint));
710 23 : assert(result);
711 23 : conn_endpoint->set_connection(connection);
712 23 : connection->set_conn_endpoint(conn_endpoint);
713 23 : return conn_endpoint;
714 23 : }
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 74 : void XmppServer::ReleaseConnectionEndpoint(XmppServerConnection *connection) {
722 74 : std::scoped_lock lock(endpoint_map_mutex_);
723 :
724 74 : if (!connection->conn_endpoint())
725 51 : return;
726 23 : assert(connection->conn_endpoint()->connection() == connection);
727 23 : connection->conn_endpoint()->reset_connection();
728 23 : connection->set_conn_endpoint(NULL);
729 74 : }
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 : }
|