Line data Source code
1 : /*
2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
3 : */
4 :
5 : #include "xmpp/xmpp_connection.h"
6 :
7 : #include <boost/date_time/posix_time/posix_time.hpp>
8 : #include <sstream>
9 :
10 : #include "base/lifetime.h"
11 : #include "base/task_annotations.h"
12 : #include "io/event_manager.h"
13 : #include "xml/xml_base.h"
14 : #include "xmpp/xmpp_client.h"
15 : #include "xmpp/xmpp_config.h"
16 : #include "xmpp/xmpp_factory.h"
17 : #include "xmpp/xmpp_log.h"
18 : #include "xmpp/xmpp_server.h"
19 : #include "xmpp/xmpp_session.h"
20 :
21 : #include "sandesh/common/vns_types.h"
22 : #include "sandesh/common/vns_constants.h"
23 : #include "sandesh/xmpp_client_server_sandesh_types.h"
24 : #include "sandesh/xmpp_message_sandesh_types.h"
25 : #include "sandesh/xmpp_server_types.h"
26 : #include "sandesh/xmpp_state_machine_sandesh_types.h"
27 : #include "sandesh/xmpp_trace_sandesh_types.h"
28 : #include "sandesh/xmpp_peer_info_types.h"
29 :
30 : using namespace std;
31 : using boost::system::error_code;
32 :
33 : const char *XmppConnection::kAuthTypeNil = "NIL";
34 : const char *XmppConnection::kAuthTypeTls = "TLS";
35 :
36 : // Maximum XMPP control message size. Typically this is less then 256 bytes. But
37 : // in scenarios where host names are quite long, we need a larger buffer size.
38 : #define XMPP_CONTROL_MESSAGE_MAX_SIZE 1024
39 :
40 92 : XmppConnection::XmppConnection(TcpServer *server,
41 92 : const XmppChannelConfig *config)
42 92 : : server_(server),
43 92 : session_(NULL),
44 92 : endpoint_(config->endpoint),
45 92 : local_endpoint_(config->local_endpoint),
46 92 : config_(NULL),
47 184 : keepalive_timer_(TimerManager::CreateTimer(
48 92 : *server->event_manager()->io_service(),
49 : "Xmpp keepalive timer",
50 : TaskScheduler::GetInstance()->GetTaskId("xmpp::StateMachine"),
51 92 : GetTaskInstance(config->ClientOnly()))),
52 92 : is_client_(config->ClientOnly()),
53 92 : log_uve_(config->logUVE),
54 92 : admin_down_(false),
55 92 : disable_read_(false),
56 92 : from_(config->FromAddr),
57 92 : to_(config->ToAddr),
58 92 : auth_enabled_(config->auth_enabled),
59 92 : dscp_value_(config->dscp_value), xmlns_(config->xmlns),
60 92 : state_machine_(XmppStaticObjectFactory::Create<XmppStateMachine>(
61 92 : this, config->ClientOnly(), config->auth_enabled, config->xmpp_hold_time)),
62 552 : mux_(XmppStaticObjectFactory::Create<XmppChannelMux>(this)) {
63 92 : ostringstream oss;
64 92 : oss << FromString() << ":" << endpoint().address().to_string();
65 92 : uve_key_str_ = oss.str();
66 92 : }
67 :
68 92 : XmppConnection::~XmppConnection() {
69 92 : StopKeepAliveTimer();
70 92 : TimerManager::DeleteTimer(keepalive_timer_);
71 92 : XMPP_UTDEBUG(XmppConnectionDelete, ToUVEKey(), XMPP_PEER_DIR_NA,
72 : "XmppConnection destructor", FromString(), ToString());
73 92 : }
74 :
75 0 : std::string XmppConnection::GetXmppAuthenticationType() const {
76 0 : if (auth_enabled_) {
77 0 : return (XmppConnection::kAuthTypeTls);
78 : } else {
79 0 : return (XmppConnection::kAuthTypeNil);
80 : }
81 : }
82 :
83 0 : void XmppConnection::SetConfig(const XmppChannelConfig *config) {
84 0 : config_ = config;
85 0 : }
86 :
87 87 : void XmppConnection::set_session(XmppSession *session) {
88 87 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
89 87 : assert(session);
90 87 : session_ = session;
91 87 : if (session_ && dscp_value_) {
92 0 : session_->SetDscpSocketOption(dscp_value_);
93 : }
94 87 : }
95 :
96 109 : void XmppConnection::clear_session() {
97 109 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
98 109 : if (!session_)
99 22 : return;
100 87 : session_->ClearConnection();
101 87 : session_ = NULL;
102 109 : }
103 :
104 25 : const XmppSession *XmppConnection::session() const {
105 25 : return session_;
106 : }
107 :
108 309 : XmppSession *XmppConnection::session() {
109 309 : return session_;
110 : }
111 :
112 0 : void XmppConnection::WriteReady() {
113 0 : boost::system::error_code ec;
114 0 : mux_->WriteReady(ec);
115 0 : }
116 :
117 24 : void XmppConnection::Shutdown() {
118 24 : ManagedDelete();
119 24 : }
120 :
121 295 : bool XmppConnection::IsDeleted() const {
122 295 : return deleter()->IsDeleted();
123 : }
124 :
125 92 : bool XmppConnection::MayDelete() const {
126 92 : return !mux_->ReceiverCount() && !mux_->RefererCount();
127 : }
128 :
129 61 : XmppSession *XmppConnection::CreateSession() {
130 61 : TcpSession *session = server_->CreateSession();
131 61 : XmppSession *xmpp_session = static_cast<XmppSession *>(session);
132 61 : xmpp_session->SetConnection(this);
133 61 : return xmpp_session;
134 : }
135 :
136 : //
137 : // Return the task instance for this XmppConnection.
138 : // Calculate from the remote IpAddress so that a restarting session uses the
139 : // same value as before.
140 : // Do not make this method virtual since it gets called from the constructor.
141 : //
142 :
143 556 : int XmppConnection::GetTaskInstance(bool is_client) const {
144 556 : if (is_client)
145 336 : return 0;
146 220 : IpAddress address = endpoint().address();
147 220 : int thread_count = TaskScheduler::GetInstance()->HardwareThreadCount();
148 220 : if (address.is_v4()) {
149 220 : return address.to_v4().to_ulong() % thread_count;
150 : } else {
151 0 : return 0;
152 : }
153 : }
154 :
155 222 : xmsm::XmState XmppConnection::GetStateMcState() const {
156 222 : const XmppStateMachine *sm = state_machine();
157 222 : assert(sm);
158 222 : return sm->StateType();
159 : }
160 :
161 141 : xmsm::XmOpenConfirmState XmppConnection::GetStateMcOpenConfirmState() const {
162 141 : const XmppStateMachine *sm = state_machine();
163 141 : assert(sm);
164 141 : return sm->OpenConfirmStateType();
165 : }
166 :
167 :
168 220 : boost::asio::ip::tcp::endpoint XmppConnection::endpoint() const {
169 220 : return endpoint_;
170 : }
171 :
172 61 : boost::asio::ip::tcp::endpoint XmppConnection::local_endpoint() const {
173 61 : return local_endpoint_;
174 : }
175 :
176 0 : string XmppConnection::endpoint_string() const {
177 0 : ostringstream oss;
178 0 : oss << endpoint_;
179 0 : return oss.str();
180 0 : }
181 :
182 0 : string XmppConnection::local_endpoint_string() const {
183 0 : ostringstream oss;
184 0 : oss << local_endpoint_;
185 0 : return oss.str();
186 0 : }
187 :
188 605 : const string &XmppConnection::FromString() const {
189 605 : return from_;
190 : }
191 :
192 665 : const string &XmppConnection::ToString() const {
193 665 : return to_;
194 : }
195 :
196 1468 : const std::string &XmppConnection::ToUVEKey() const {
197 1468 : return uve_key_str_;
198 : }
199 :
200 0 : static void XMPPPeerInfoSend(XmppPeerInfoData &peer_info) {
201 0 : assert(!peer_info.get_name().empty());
202 0 : XMPPPeerInfo::Send(peer_info);
203 0 : }
204 :
205 86 : void XmppConnection::SetTo(const string &to) {
206 86 : if ((to_.size() == 0) && (to.size() != 0)) {
207 55 : to_ = to;
208 55 : if (!logUVE()) return;
209 0 : XmppPeerInfoData peer_info;
210 0 : peer_info.set_name(ToUVEKey());
211 0 : peer_info.set_identifier(to_);
212 0 : XMPPPeerInfoSend(peer_info);
213 0 : }
214 : }
215 :
216 0 : void XmppConnection::SetAdminDown(bool toggle) {
217 : // TODO: generate state machine event.
218 0 : admin_down_ = toggle;
219 0 : }
220 :
221 14 : bool XmppConnection::AcceptSession(XmppSession *session) {
222 14 : session->SetConnection(this);
223 14 : return state_machine_->PassiveOpen(session);
224 : }
225 :
226 22 : bool XmppConnection::Send(const uint8_t *data, size_t size,
227 : const string *msg_str) {
228 22 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
229 22 : if (session_ == NULL) {
230 0 : return false;
231 : }
232 :
233 22 : TcpSession::Endpoint endpoint = session_->remote_endpoint();
234 22 : const string &endpoint_addr_str = session_->remote_addr_string();
235 22 : string str;
236 22 : if (!msg_str) {
237 22 : str.append(reinterpret_cast<const char *>(data), size);
238 22 : msg_str = &str;
239 : }
240 :
241 22 : if (!(mux_ &&
242 22 : (mux_->TxMessageTrace(endpoint_addr_str, endpoint.port(),
243 22 : size, *msg_str, NULL)))) {
244 22 : XMPP_MESSAGE_TRACE(XmppTxStream,
245 : endpoint_addr_str, endpoint.port(), size, *msg_str);
246 : }
247 :
248 22 : stats_[1].update++;
249 : size_t sent;
250 22 : return session_->Send(data, size, &sent);
251 22 : }
252 :
253 0 : int XmppConnection::SetDscpValue(uint8_t value) {
254 0 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
255 0 : dscp_value_ = value;
256 0 : if (!session_) {
257 0 : return 0;
258 : }
259 0 : return session_->SetDscpSocketOption(value);
260 0 : }
261 :
262 60 : bool XmppConnection::SendOpen(XmppSession *session) {
263 60 : if (!session) return false;
264 60 : XmppProto::XmppStanza::XmppStreamMessage openstream;
265 60 : openstream.strmtype = XmppStanza::XmppStreamMessage::INIT_STREAM_HEADER;
266 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
267 60 : int len = XmppProto::EncodeStream(openstream, to_, from_, xmlns_, data,
268 : sizeof(data));
269 60 : if (len <= 0) {
270 0 : inc_open_fail();
271 0 : return false;
272 : } else {
273 60 : XMPP_UTDEBUG(XmppOpen, ToUVEKey(), XMPP_PEER_DIR_OUT, len, from_, to_,
274 : xmlns_);
275 60 : session->Send(data, len, NULL);
276 60 : stats_[1].open++;
277 60 : return true;
278 : }
279 60 : }
280 :
281 37 : bool XmppConnection::SendOpenConfirm(XmppSession *session) {
282 37 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
283 37 : if (!session_) return false;
284 37 : XmppStanza::XmppStreamMessage openstream;
285 37 : openstream.strmtype = XmppStanza::XmppStreamMessage::INIT_STREAM_HEADER_RESP;
286 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
287 37 : int len = XmppProto::EncodeStream(openstream, to_, from_, xmlns_, data,
288 : sizeof(data));
289 37 : if (len <= 0) {
290 0 : inc_open_fail();
291 0 : return false;
292 : } else {
293 37 : XMPP_UTDEBUG(XmppOpenConfirm, ToUVEKey(), XMPP_PEER_DIR_OUT, len,
294 : from_, to_);
295 37 : session_->Send(data, len, NULL);
296 37 : stats_[1].open++;
297 37 : return true;
298 : }
299 37 : }
300 :
301 14 : bool XmppConnection::SendStreamFeatureRequest(XmppSession *session) {
302 14 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
303 14 : if (!session_) return false;
304 14 : XmppStanza::XmppStreamMessage featurestream;
305 14 : featurestream.strmtype = XmppStanza::XmppStreamMessage::FEATURE_TLS;
306 14 : featurestream.strmtlstype = XmppStanza::XmppStreamMessage::TLS_FEATURE_REQUEST;
307 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
308 14 : int len = XmppProto::EncodeStream(featurestream, to_, from_, xmlns_, data,
309 : sizeof(data));
310 14 : if (len <= 0) {
311 0 : inc_stream_feature_fail();
312 0 : return false;
313 : } else {
314 14 : XMPP_UTDEBUG(XmppControlMessage, ToUVEKey(), XMPP_PEER_DIR_OUT,
315 : "Send Stream Feature Request", len, from_, to_);
316 14 : session_->Send(data, len, NULL);
317 : //stats_[1].open++;
318 14 : return true;
319 : }
320 14 : }
321 :
322 15 : bool XmppConnection::SendStartTls(XmppSession *session) {
323 15 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
324 15 : if (!session_) return false;
325 15 : XmppStanza::XmppStreamMessage stream;
326 15 : stream.strmtype = XmppStanza::XmppStreamMessage::FEATURE_TLS;
327 15 : stream.strmtlstype = XmppStanza::XmppStreamMessage::TLS_START;
328 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
329 15 : int len = XmppProto::EncodeStream(stream, to_, from_, xmlns_, data,
330 : sizeof(data));
331 15 : if (len <= 0) {
332 0 : inc_stream_feature_fail();
333 0 : return false;
334 : } else {
335 15 : XMPP_UTDEBUG(XmppControlMessage, ToUVEKey(), XMPP_PEER_DIR_OUT,
336 : "Send Start Tls", len, from_, to_);
337 15 : session_->Send(data, len, NULL);
338 : //stats_[1].open++;
339 15 : return true;
340 : }
341 15 : }
342 :
343 10 : bool XmppConnection::SendProceedTls(XmppSession *session) {
344 10 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
345 10 : if (!session_) return false;
346 10 : XmppStanza::XmppStreamMessage stream;
347 10 : stream.strmtype = XmppStanza::XmppStreamMessage::FEATURE_TLS;
348 10 : stream.strmtlstype = XmppStanza::XmppStreamMessage::TLS_PROCEED;
349 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
350 10 : int len = XmppProto::EncodeStream(stream, to_, from_, xmlns_, data,
351 : sizeof(data));
352 10 : if (len <= 0) {
353 0 : inc_stream_feature_fail();
354 0 : return false;
355 : } else {
356 10 : XMPP_UTDEBUG(XmppControlMessage, ToUVEKey(), XMPP_PEER_DIR_OUT,
357 : "Send Proceed Tls", len, from_, to_);
358 10 : session_->Send(data, len, NULL);
359 : //stats_[1].open++;
360 10 : return true;
361 : }
362 10 : }
363 :
364 2 : void XmppConnection::SendClose(XmppSession *session) {
365 2 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
366 2 : if (!session_) return;
367 1 : string str("</stream:stream>");
368 : uint8_t data[64];
369 1 : memcpy(data, str.data(), str.size());
370 1 : XMPP_UTDEBUG(XmppClose, ToUVEKey(), XMPP_PEER_DIR_OUT, str.size(), from_,
371 : to_);
372 1 : session_->Send(data, str.size(), NULL);
373 1 : stats_[1].close++;
374 2 : }
375 :
376 2 : void XmppConnection::ProcessSslHandShakeResponse(SslSessionPtr session,
377 : const boost::system::error_code& error) {
378 2 : if (!state_machine())
379 0 : return;
380 :
381 2 : if (error) {
382 0 : inc_handshake_failure();
383 :
384 0 : if (error.category() == boost::asio::error::get_ssl_category()) {
385 0 : string err = error.message();
386 0 : err = string(" (")
387 0 : +boost::lexical_cast<string>(ERR_GET_LIB(error.value()))+","
388 0 : +boost::lexical_cast<string>(ERR_GET_REASON(error.value()))+") ";
389 :
390 : char buf[128];
391 0 : ::ERR_error_string_n(error.value(), buf, sizeof(buf));
392 0 : err += buf;
393 0 : XMPP_ALERT(XmppSslHandShakeFailure, ToUVEKey(), XMPP_PEER_DIR_IN,
394 : "failure", err);
395 0 : }
396 :
397 0 : state_machine()->OnEvent(session.get(), xmsm::EvTLSHANDSHAKE_FAILURE);
398 :
399 : } else {
400 2 : XMPP_DEBUG(XmppSslHandShakeMessage, session->ToUVEKey(),
401 : XMPP_PEER_DIR_IN, "success", "");
402 2 : state_machine()->OnEvent(session.get(), xmsm::EvTLSHANDSHAKE_SUCCESS);
403 : }
404 : }
405 :
406 0 : void XmppConnection::LogMsg(std::string msg) {
407 0 : log4cplus::Logger logger = log4cplus::Logger::getRoot();
408 0 : LOG4CPLUS_DEBUG(logger, msg << ToString() << " " <<
409 : local_endpoint_.address() << ":" << local_endpoint_.port() << "::" <<
410 : endpoint_.address() << ":" << endpoint_.port());
411 0 : }
412 :
413 30 : void XmppConnection::LogKeepAliveSend() {
414 : static bool init_ = false;
415 : static bool log_ = false;
416 :
417 30 : if (!init_) {
418 5 : char *str = getenv("XMPP_ASSERT_ON_HOLD_TIMEOUT");
419 5 : if (str && strtoul(str, NULL, 0) != 0) log_ = true;
420 5 : init_ = true;
421 : }
422 :
423 30 : if (log_) LogMsg("SEND KEEPALIVE: ");
424 30 : }
425 :
426 30 : void XmppConnection::SendKeepAlive() {
427 30 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
428 30 : if (!session_) return;
429 30 : XmppStanza::XmppMessage msg(XmppStanza::WHITESPACE_MESSAGE_STANZA);
430 : uint8_t data[XMPP_CONTROL_MESSAGE_MAX_SIZE];
431 30 : int len = XmppProto::EncodeStream(msg, data, sizeof(data));
432 30 : assert(len > 0);
433 30 : session_->Send(data, len, NULL);
434 30 : stats_[1].keepalive++;
435 30 : LogKeepAliveSend();
436 30 : }
437 :
438 0 : bool XmppConnection::KeepAliveTimerExpired() {
439 0 : if (state_machine_->get_state() != xmsm::ESTABLISHED)
440 0 : return false;
441 :
442 : // TODO: check timestamp of last received packet.
443 0 : SendKeepAlive();
444 :
445 : //
446 : // Start the timer again, by returning true
447 : //
448 0 : return true;
449 : }
450 :
451 0 : void XmppConnection::KeepaliveTimerErrorHanlder(string error_name,
452 : string error_message) {
453 0 : XMPP_WARNING(XmppKeepaliveTimeError, ToUVEKey(), XMPP_PEER_DIR_NA,
454 : error_name, error_message);
455 0 : }
456 :
457 53 : void XmppConnection::StartKeepAliveTimer() {
458 53 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
459 53 : if (!session_)
460 0 : return;
461 :
462 53 : int holdtime_msecs = state_machine_->hold_time_msecs();
463 53 : if (holdtime_msecs <= 0)
464 0 : return;
465 :
466 53 : keepalive_timer_->Start(holdtime_msecs / 3,
467 : boost::bind(&XmppConnection::KeepAliveTimerExpired, this),
468 : boost::bind(&XmppConnection::KeepaliveTimerErrorHanlder, this, _1, _2));
469 53 : }
470 :
471 141 : void XmppConnection::StopKeepAliveTimer() {
472 141 : tbb::spin_mutex::scoped_lock lock(spin_mutex_);
473 141 : keepalive_timer_->Cancel();
474 141 : }
475 :
476 0 : void XmppConnection::UpdateKeepAliveTimer(uint8_t time_out) {
477 0 : state_machine_->set_hold_time(time_out);
478 0 : StopKeepAliveTimer();
479 0 : StartKeepAliveTimer();
480 0 : }
481 :
482 381 : XmppStateMachine *XmppConnection::state_machine() {
483 381 : return state_machine_.get();
484 : }
485 :
486 363 : const XmppStateMachine *XmppConnection::state_machine() const {
487 363 : return state_machine_.get();
488 : }
489 :
490 0 : const XmppChannelMux *XmppConnection::channel_mux() const {
491 0 : return mux_.get();
492 : }
493 :
494 84 : void XmppConnection::IncProtoStats(unsigned int type) {
495 84 : switch (type) {
496 31 : case XmppStanza::STREAM_HEADER:
497 31 : stats_[0].open++;
498 31 : break;
499 30 : case XmppStanza::WHITESPACE_MESSAGE_STANZA:
500 30 : stats_[0].keepalive++;
501 30 : break;
502 23 : case XmppStanza::IQ_STANZA:
503 23 : stats_[0].update++;
504 23 : break;
505 0 : case XmppStanza::MESSAGE_STANZA:
506 0 : stats_[0].update++;
507 0 : break;
508 : }
509 84 : }
510 :
511 5 : void XmppConnection::inc_connect_error() {
512 5 : error_stats_.connect_error++;
513 5 : }
514 :
515 28 : void XmppConnection::inc_session_close() {
516 28 : error_stats_.session_close++;
517 28 : }
518 :
519 0 : void XmppConnection::inc_open_fail() {
520 0 : error_stats_.open_fail++;
521 0 : }
522 :
523 0 : void XmppConnection::inc_stream_feature_fail() {
524 0 : error_stats_.stream_feature_fail++;
525 0 : }
526 :
527 0 : void XmppConnection::inc_handshake_failure() {
528 0 : error_stats_.handshake_fail++;
529 0 : }
530 :
531 0 : size_t XmppConnection::get_connect_error() {
532 0 : return error_stats_.connect_error;
533 : }
534 :
535 0 : size_t XmppConnection::get_session_close() {
536 0 : return error_stats_.session_close;
537 : }
538 :
539 0 : size_t XmppConnection::get_open_fail() {
540 0 : return error_stats_.open_fail;
541 : }
542 :
543 0 : size_t XmppConnection::get_stream_feature_fail() {
544 0 : return error_stats_.stream_feature_fail;
545 : }
546 :
547 0 : size_t XmppConnection::get_handshake_failure() {
548 0 : return error_stats_.handshake_fail;
549 : }
550 :
551 0 : size_t XmppConnection::get_sm_connect_attempts() {
552 0 : return state_machine_->get_connect_attempts();
553 : }
554 :
555 0 : size_t XmppConnection::get_sm_keepalive_count() {
556 0 : return state_machine_->get_keepalive_count();
557 : }
558 :
559 84 : void XmppConnection::ReceiveMsg(XmppSession *session, const string &msg) {
560 84 : XmppStanza::XmppMessage *minfo = XmppDecode(msg);
561 :
562 84 : if (minfo) {
563 82 : session->IncStats((unsigned int)minfo->type, msg.size());
564 82 : if (minfo->type != XmppStanza::WHITESPACE_MESSAGE_STANZA) {
565 52 : if (!(mux_ &&
566 104 : (mux_->RxMessageTrace(session->
567 104 : remote_endpoint().address().to_string(),
568 104 : session->remote_endpoint().port(),
569 104 : msg.size(), msg, minfo)))) {
570 52 : XMPP_MESSAGE_TRACE(XmppRxStream,
571 : session->
572 : remote_endpoint().address().to_string(),
573 : session->
574 : remote_endpoint().port(), msg.size(), msg);
575 : }
576 : }
577 82 : IncProtoStats((unsigned int)minfo->type);
578 82 : state_machine_->OnMessage(session, minfo);
579 2 : } else if ((minfo = last_msg_.get()) != NULL) {
580 2 : session->IncStats((unsigned int)minfo->type, msg.size());
581 2 : IncProtoStats((unsigned int)minfo->type);
582 : } else {
583 0 : session->IncStats(XmppStanza::INVALID, msg.size());
584 0 : XMPP_MESSAGE_TRACE(XmppRxStreamInvalid,
585 : session->remote_endpoint().address().to_string(),
586 : session->remote_endpoint().port(), msg.size(), msg);
587 : }
588 84 : return;
589 : }
590 :
591 84 : XmppStanza::XmppMessage *XmppConnection::XmppDecode(const string &msg) {
592 84 : unique_ptr<XmppStanza::XmppMessage> minfo(XmppProto::Decode(this, msg));
593 84 : if (minfo.get() == NULL) {
594 0 : XMPP_INFO(XmppSessionDelete, ToUVEKey(), XMPP_PEER_DIR_IN, "Server",
595 : FromString(), ToString());
596 0 : Clear();
597 0 : return NULL;
598 : }
599 :
600 84 : if (minfo->type == XmppStanza::IQ_STANZA) {
601 : const XmppStanza::XmppMessageIq *iq =
602 23 : static_cast<const XmppStanza::XmppMessageIq *>(minfo.get());
603 :
604 :
605 23 : if (iq->action.compare("publish") == 0) {
606 2 : last_msg_.reset(minfo.release());
607 2 : return NULL;
608 : }
609 :
610 21 : if (iq->action.compare("collection") == 0) {
611 2 : if (last_msg_.get() != NULL) {
612 : XmppStanza::XmppMessageIq *last_iq =
613 2 : static_cast<XmppStanza::XmppMessageIq *>(last_msg_.get());
614 :
615 2 : if (last_iq->node.compare(iq->as_node) == 0) {
616 2 : XmlBase *impl = last_iq->dom.get();
617 2 : impl->ReadNode("publish");
618 2 : impl->ModifyAttribute("node", iq->node);
619 2 : last_iq->node = impl->ReadAttrib("node");
620 2 : last_iq->is_as_node = iq->is_as_node;
621 : //Save the complete ass/dissociate node
622 2 : last_iq->as_node = iq->as_node;
623 : } else {
624 0 : XMPP_WARNING(XmppIqMessageInvalid, ToUVEKey(),
625 : XMPP_PEER_DIR_IN);
626 0 : goto error;
627 : }
628 : } else {
629 0 : XMPP_ERROR(XmppIqCollectionError, ToUVEKey(), XMPP_PEER_DIR_IN);
630 0 : goto error;
631 : }
632 : // iq message merged with collection info
633 2 : return last_msg_.release();
634 : }
635 : }
636 80 : return minfo.release();
637 :
638 0 : error:
639 0 : last_msg_.reset();
640 0 : return NULL;
641 84 : }
642 :
643 4 : int XmppConnection::ProcessXmppChatMessage(
644 : const XmppStanza::XmppChatMessage *msg) {
645 4 : mux_->ProcessXmppMessage(msg);
646 4 : return 0;
647 : }
648 :
649 21 : int XmppConnection::ProcessXmppIqMessage(const XmppStanza::XmppMessage *msg) {
650 21 : mux_->ProcessXmppMessage(msg);
651 21 : return 0;
652 : }
653 :
654 : class XmppServerConnection::DeleteActor : public LifetimeActor {
655 : public:
656 37 : DeleteActor(XmppServer *server, XmppServerConnection *parent)
657 37 : : LifetimeActor(server->lifetime_manager()),
658 37 : server_(server), parent_(parent) {
659 37 : }
660 :
661 37 : virtual bool MayDelete() const {
662 37 : return (!parent_->on_work_queue() && parent_->MayDelete());
663 : }
664 :
665 37 : virtual void Shutdown() {
666 37 : CHECK_CONCURRENCY("bgp::Config");
667 :
668 : // If the connection is still on the WorkQueue, simply add it to the
669 : // ConnectionSet. It won't be on the ConnectionMap.
670 37 : if (parent_->on_work_queue()) {
671 0 : server_->InsertDeletedConnection(parent_);
672 : }
673 :
674 : // If the connection was rejected as duplicate, it will already be in
675 : // the ConnectionSet. Non-duplicate connections need to be moved from
676 : // from the ConnectionMap into the ConnectionSet. We add it to the
677 : // ConnectionSet and then remove it from ConnectionMap to ensure that
678 : // the XmppServer connection count doesn't temporarily become 0. This
679 : // is friendly to tests that wait for the XmppServer connection count
680 : // to become 0.
681 : //
682 : // Breaking association with the XmppConnectionEndpoint here allows
683 : // a new XmppServerConnection with the same Endpoint to come up. We
684 : // may end up leaking memory if current XmppServerConnection doesn't
685 : // get cleaned up completely, but we at least prevent the other end
686 : // from getting stuck forever.
687 37 : else if (!parent_->duplicate()) {
688 37 : server_->InsertDeletedConnection(parent_);
689 37 : server_->RemoveConnection(parent_);
690 37 : server_->ReleaseConnectionEndpoint(parent_);
691 : }
692 :
693 37 : if (parent_->state_machine()) {
694 37 : parent_->state_machine()->Clear();
695 : }
696 :
697 37 : XmppSession *session = NULL;
698 37 : if (parent_->state_machine()) {
699 37 : session = parent_->state_machine()->session();
700 37 : parent_->state_machine()->clear_session();
701 : }
702 37 : if (session) {
703 9 : server_->DeleteSession(session);
704 : }
705 37 : }
706 :
707 37 : virtual void Destroy() {
708 37 : delete parent_;
709 37 : }
710 :
711 : private:
712 : XmppServer *server_;
713 : XmppServerConnection *parent_;
714 : };
715 :
716 37 : XmppServerConnection::XmppServerConnection(XmppServer *server,
717 37 : const XmppChannelConfig *config)
718 : : XmppConnection(server, config),
719 37 : duplicate_(false),
720 37 : on_work_queue_(false),
721 37 : conn_endpoint_(NULL),
722 37 : deleter_(new DeleteActor(server, this)),
723 74 : server_delete_ref_(this, server->deleter()) {
724 37 : assert(!config->ClientOnly());
725 37 : XMPP_INFO(XmppConnectionCreate, ToUVEKey(), XMPP_PEER_DIR_IN,
726 : "Server", FromString(), ToString());
727 37 : }
728 :
729 61 : XmppServerConnection::~XmppServerConnection() {
730 37 : CHECK_CONCURRENCY("bgp::Config");
731 :
732 37 : XMPP_INFO(XmppConnectionDelete, ToUVEKey(), XMPP_PEER_DIR_NA, "Server",
733 : FromString(), ToString());
734 37 : server()->RemoveDeletedConnection(this);
735 61 : }
736 :
737 53 : void XmppServerConnection::ManagedDelete() {
738 53 : XMPP_UTDEBUG(XmppConnectionDelete, ToUVEKey(), XMPP_PEER_DIR_NA,
739 : "Managed server connection delete", FromString(), ToString());
740 53 : deleter_->Delete();
741 53 : }
742 :
743 17 : void XmppServerConnection::RetryDelete() {
744 17 : if (!deleter()->IsDeleted())
745 17 : return;
746 0 : deleter()->RetryDelete();
747 : }
748 :
749 0 : LifetimeManager *XmppServerConnection::lifetime_manager() {
750 0 : return server()->lifetime_manager();
751 : }
752 :
753 54 : XmppServer *XmppServerConnection::server() {
754 54 : return static_cast<XmppServer *>(server_);
755 : }
756 :
757 65 : LifetimeActor *XmppServerConnection::deleter() {
758 65 : return deleter_.get();
759 : }
760 :
761 210 : const LifetimeActor *XmppServerConnection::deleter() const {
762 210 : return deleter_.get();
763 : }
764 :
765 2 : void XmppServerConnection::set_close_reason(const string &close_reason) {
766 2 : if (conn_endpoint_)
767 0 : conn_endpoint_->set_close_reason(close_reason);
768 :
769 2 : if (!logUVE())
770 2 : return;
771 :
772 0 : XmppPeerInfoData peer_info;
773 0 : peer_info.set_name(ToUVEKey());
774 0 : peer_info.set_close_reason(close_reason);
775 0 : XMPPPeerInfoSend(peer_info);
776 0 : }
777 :
778 0 : uint32_t XmppServerConnection::flap_count() const {
779 0 : return conn_endpoint_ ? conn_endpoint_->flap_count() : 0;
780 : }
781 :
782 19 : void XmppServerConnection::increment_flap_count() {
783 19 : XmppConnectionEndpoint *conn_endpoint = conn_endpoint_;
784 19 : if (!conn_endpoint)
785 0 : conn_endpoint = server()->FindConnectionEndpoint(ToString());
786 19 : if (!conn_endpoint)
787 19 : return;
788 19 : conn_endpoint->increment_flap_count();
789 :
790 19 : if (!logUVE())
791 19 : return;
792 :
793 0 : XmppPeerInfoData peer_info;
794 0 : peer_info.set_name(ToUVEKey());
795 0 : PeerFlapInfo flap_info;
796 0 : flap_info.set_flap_count(conn_endpoint->flap_count());
797 0 : flap_info.set_flap_time(conn_endpoint->last_flap());
798 0 : peer_info.set_flap_info(flap_info);
799 0 : XMPPPeerInfoSend(peer_info);
800 0 : }
801 :
802 0 : const std::string XmppServerConnection::last_flap_at() const {
803 0 : return conn_endpoint_ ? conn_endpoint_->last_flap_at() : "";
804 : }
805 :
806 0 : void XmppServerConnection::FillShowInfo(
807 : ShowXmppConnection *show_connection) const {
808 0 : show_connection->set_name(ToString());
809 0 : show_connection->set_deleted(IsDeleted());
810 0 : show_connection->set_remote_endpoint(endpoint_string());
811 0 : show_connection->set_local_endpoint(local_endpoint_string());
812 0 : show_connection->set_state(StateName());
813 0 : show_connection->set_last_event(LastEvent());
814 0 : show_connection->set_last_state(LastStateName());
815 0 : show_connection->set_last_state_at(LastStateChangeAt());
816 0 : show_connection->set_receivers(channel_mux()->GetReceiverList());
817 0 : show_connection->set_server_auth_type(GetXmppAuthenticationType());
818 0 : show_connection->set_dscp_value(dscp_value());
819 0 : }
820 :
821 : class XmppClientConnection::DeleteActor : public LifetimeActor {
822 : public:
823 55 : DeleteActor(XmppClient *client, XmppClientConnection *parent)
824 55 : : LifetimeActor(client->lifetime_manager()),
825 55 : client_(client), parent_(parent) {
826 55 : }
827 :
828 55 : virtual bool MayDelete() const {
829 55 : return parent_->MayDelete();
830 : }
831 :
832 55 : virtual void Shutdown() {
833 55 : if (parent_->session()) {
834 25 : client_->NotifyConnectionEvent(parent_->ChannelMux(),
835 : xmps::NOT_READY);
836 : }
837 :
838 55 : XmppSession *session = NULL;
839 55 : if (parent_->state_machine()) {
840 55 : session = parent_->state_machine()->session();
841 55 : parent_->state_machine()->clear_session();
842 : }
843 55 : if (session) {
844 30 : client_->DeleteSession(session);
845 : }
846 55 : }
847 :
848 55 : virtual void Destroy() {
849 55 : delete parent_;
850 55 : }
851 :
852 : private:
853 : XmppClient *client_;
854 : XmppClientConnection *parent_;
855 : };
856 :
857 55 : XmppClientConnection::XmppClientConnection(XmppClient *server,
858 55 : const XmppChannelConfig *config)
859 : : XmppConnection(server, config),
860 55 : flap_count_(0),
861 55 : deleter_(new DeleteActor(server, this)),
862 110 : server_delete_ref_(this, server->deleter()) {
863 55 : assert(config->ClientOnly());
864 55 : XMPP_UTDEBUG(XmppConnectionCreate, ToUVEKey(), XMPP_PEER_DIR_NA, "Client",
865 : FromString(), ToString());
866 55 : }
867 :
868 83 : XmppClientConnection::~XmppClientConnection() {
869 55 : CHECK_CONCURRENCY("bgp::Config");
870 :
871 55 : XMPP_INFO(XmppConnectionDelete, ToUVEKey(), XMPP_PEER_DIR_NA,
872 : "Client", FromString(), ToString());
873 55 : server()->RemoveConnection(this);
874 83 : }
875 :
876 55 : void XmppClientConnection::ManagedDelete() {
877 55 : XMPP_UTDEBUG(XmppConnectionDelete, ToUVEKey(), XMPP_PEER_DIR_NA,
878 : "Managed Client Delete", FromString(), ToString());
879 55 : deleter_->Delete();
880 55 : }
881 :
882 3 : void XmppClientConnection::RetryDelete() {
883 3 : if (!deleter()->IsDeleted())
884 3 : return;
885 0 : deleter()->RetryDelete();
886 : }
887 :
888 0 : LifetimeManager *XmppClientConnection::lifetime_manager() {
889 0 : return server()->lifetime_manager();
890 : }
891 :
892 55 : XmppClient *XmppClientConnection::server() {
893 55 : return static_cast<XmppClient *>(server_);
894 : }
895 :
896 3 : LifetimeActor *XmppClientConnection::deleter() {
897 3 : return deleter_.get();
898 : }
899 :
900 85 : const LifetimeActor *XmppClientConnection::deleter() const {
901 85 : return deleter_.get();
902 : }
903 :
904 10 : void XmppClientConnection::set_close_reason(const string &close_reason) {
905 10 : close_reason_ = close_reason;
906 10 : if (!logUVE())
907 10 : return;
908 :
909 0 : XmppPeerInfoData peer_info;
910 0 : peer_info.set_name(ToUVEKey());
911 0 : peer_info.set_close_reason(close_reason_);
912 0 : XMPPPeerInfoSend(peer_info);
913 0 : }
914 :
915 0 : uint32_t XmppClientConnection::flap_count() const {
916 0 : return flap_count_;
917 : }
918 :
919 10 : void XmppClientConnection::increment_flap_count() {
920 10 : flap_count_++;
921 10 : last_flap_ = UTCTimestampUsec();
922 :
923 10 : if (!logUVE())
924 10 : return;
925 :
926 0 : XmppPeerInfoData peer_info;
927 0 : peer_info.set_name(ToUVEKey());
928 0 : PeerFlapInfo flap_info;
929 0 : flap_info.set_flap_count(flap_count_);
930 0 : flap_info.set_flap_time(last_flap_);
931 0 : peer_info.set_flap_info(flap_info);
932 0 : XMPPPeerInfoSend(peer_info);
933 0 : }
934 :
935 0 : const std::string XmppClientConnection::last_flap_at() const {
936 0 : return last_flap_ ? integerToString(UTCUsecToPTime(last_flap_)) : "";
937 : }
938 :
939 23 : XmppConnectionEndpoint::XmppConnectionEndpoint(const string &client)
940 23 : : client_(client), flap_count_(0), last_flap_(0), connection_(NULL) {
941 23 : }
942 :
943 0 : void XmppConnectionEndpoint::set_close_reason(const string &close_reason) {
944 0 : close_reason_ = close_reason;
945 0 : }
946 :
947 0 : uint32_t XmppConnectionEndpoint::flap_count() const {
948 0 : return flap_count_;
949 : }
950 :
951 19 : void XmppConnectionEndpoint::increment_flap_count() {
952 19 : flap_count_++;
953 19 : last_flap_ = UTCTimestampUsec();
954 19 : }
955 :
956 0 : uint64_t XmppConnectionEndpoint::last_flap() const {
957 0 : return last_flap_;
958 : }
959 :
960 0 : const std::string XmppConnectionEndpoint::last_flap_at() const {
961 0 : return last_flap_ ? integerToString(UTCUsecToPTime(last_flap_)) : "";
962 : }
963 :
964 25 : XmppConnection *XmppConnectionEndpoint::connection() {
965 25 : return connection_;
966 : }
967 :
968 0 : const XmppConnection *XmppConnectionEndpoint::connection() const {
969 0 : return connection_;
970 : }
971 :
972 23 : void XmppConnectionEndpoint::set_connection(XmppConnection *connection) {
973 23 : assert(!connection_);
974 23 : connection_ = connection;
975 23 : }
976 :
977 23 : void XmppConnectionEndpoint::reset_connection() {
978 23 : assert(connection_);
979 23 : connection_ = NULL;
980 23 : }
981 :
982 : // Swap relavent contents between two XmppConnection objects.
983 0 : void XmppConnection::SwapContents(XmppConnection *other) {
984 0 : assert(!IsClient());
985 0 : assert(!other->IsClient());
986 : // Update the ConnectionMap in the server as the endpoints are the keys.
987 0 : XmppServer *server = dynamic_cast<XmppServerConnection *>(this)->server();
988 0 : server->SwapXmppConnectionMapEntries(this, other);
989 : // Swap all other connection related information.
990 0 : swap(local_endpoint_, other->local_endpoint_);
991 0 : stats_[0].swap(other->stats_[0]);
992 0 : stats_[1].swap(other->stats_[1]);
993 0 : error_stats_.swap(other->error_stats_);
994 0 : swap(last_msg_, other->last_msg_);
995 0 : swap(to_, other->to_);
996 0 : swap(from_, other->from_);
997 0 : swap(xmlns_, other->xmlns_);
998 0 : swap(dscp_value_, other->dscp_value_);
999 0 : swap(disable_read_, other->disable_read_);
1000 0 : }
|