Line data Source code
1 : /* 2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved. 3 : */ 4 : 5 : #include "base/regex.h" 6 : #include "xmpp/xmpp_session.h" 7 : 8 : #include "xmpp/xmpp_connection.h" 9 : #include "xmpp/xmpp_log.h" 10 : #include "xmpp/xmpp_proto.h" 11 : #include "xmpp/xmpp_server.h" 12 : #include "xmpp/xmpp_state_machine.h" 13 : 14 : #include "sandesh/sandesh_trace.h" 15 : #include "sandesh/xmpp_trace_sandesh_types.h" 16 : 17 : using namespace std; 18 : using contrail::regex; 19 : using contrail::regex_match; 20 : using contrail::regex_search; 21 : 22 : using boost::asio::mutable_buffer; 23 : 24 : const regex XmppSession::patt_(rXMPP_MESSAGE); 25 : const regex XmppSession::stream_patt_(rXMPP_STREAM_START); 26 : const regex XmppSession::stream_res_end_(rXMPP_STREAM_END); 27 : const regex XmppSession::whitespace_(sXMPP_WHITESPACE); 28 : const regex XmppSession::stream_features_patt_(rXMPP_STREAM_FEATURES); 29 : const regex XmppSession::starttls_patt_(rXMPP_STREAM_STARTTLS); 30 : const regex XmppSession::proceed_patt_(rXMPP_STREAM_PROCEED); 31 : const regex XmppSession::end_patt_(rXMPP_STREAM_STANZA_END); 32 : 33 96 : XmppSession::XmppSession(XmppConnectionManager *manager, SslSocket *socket, 34 96 : bool async_ready) 35 : : SslSession(manager, socket, async_ready), 36 96 : manager_(manager), 37 96 : connection_(NULL), 38 96 : tag_known_(0), 39 96 : task_instance_(-1), 40 96 : stats_(XmppStanza::RESERVED_STANZA, XmppSession::StatsPair(0, 0)), 41 96 : keepalive_probes_(kSessionKeepaliveProbes) { 42 96 : buf_.reserve(kMaxMessageSize); 43 96 : offset_ = buf_.begin(); 44 96 : stream_open_matched_ = false; 45 96 : } 46 : 47 191 : XmppSession::~XmppSession() { 48 96 : set_observer(NULL); 49 96 : connection_ = NULL; 50 191 : } 51 : 52 95 : void XmppSession::SetConnection(XmppConnection *connection) { 53 95 : assert(connection); 54 95 : connection_ = connection; 55 95 : task_instance_ = connection_->GetTaskInstance(); 56 95 : } 57 : 58 : // 59 : // Dissociate the connection from the this XmppSession. 60 : // Do not invalidate the task_instance since it can be used to spawn an 61 : // io::ReaderTask while this method is being executed. 62 : // 63 182 : void XmppSession::ClearConnection() { 64 182 : connection_ = NULL; 65 182 : } 66 : 67 : // 68 : // Concurrency: called in the context of bgp::Config task. 69 : // 70 : // Process write ready callback. 71 : // 72 0 : void XmppSession::ProcessWriteReady() { 73 0 : if (!connection_) 74 0 : return; 75 0 : connection_->WriteReady(); 76 : } 77 : 78 : // 79 : // Concurrency: called in the context of io thread. 80 : // 81 : // Handle write ready callback. 82 : // 83 : // Enqueue session to the XmppConnectionManager. The session is added to a 84 : // WorkQueue gets processed in the context of bgp::Config task. Doing this 85 : // ensures that we don't access the XmppConnection while the XmppConnection 86 : // is trying to clear our back pointer to it. 87 : // 88 : // We can ignore any errors since the StateMachine will get informed of the 89 : // TcpSession close independently and react to it. 90 : // 91 0 : void XmppSession::WriteReady(const boost::system::error_code &error) { 92 0 : if (error) 93 0 : return; 94 0 : manager_->EnqueueSession(this); 95 : } 96 : 97 30 : XmppSession::StatsPair XmppSession::Stats(unsigned int type) const { 98 30 : assert (type < (unsigned int)XmppStanza::RESERVED_STANZA); 99 30 : return stats_[type]; 100 : } 101 : 102 83 : void XmppSession::IncStats(unsigned int type, uint64_t bytes) { 103 83 : assert (type < (unsigned int)XmppStanza::RESERVED_STANZA); 104 83 : stats_[type].first++; 105 83 : stats_[type].second += bytes; 106 83 : } 107 : 108 57 : boost::system::error_code XmppSession::EnableTcpKeepalive(int hold_time) { 109 57 : char *keepalive_time_str = getenv("TCP_KEEPALIVE_SECONDS"); 110 57 : if (keepalive_time_str) { 111 0 : hold_time = strtoul(keepalive_time_str, NULL, 0) * 3; 112 0 : if (!hold_time) 113 0 : return boost::system::error_code(); 114 : } 115 : 116 57 : if (hold_time <= 9) { 117 0 : hold_time = 9; // min hold-time in secs. 118 : } 119 57 : hold_time = ((hold_time > 18)? hold_time/2 : hold_time); 120 57 : keepalive_idle_time_ = hold_time/3; 121 57 : keepalive_interval_ = 122 57 : ((hold_time - keepalive_idle_time_)/keepalive_probes_); 123 57 : tcp_user_timeout_ = (hold_time * 1000); // msec 124 : 125 57 : return (SetSocketKeepaliveOptions(keepalive_idle_time_, 126 : keepalive_interval_, 127 : keepalive_probes_, 128 57 : tcp_user_timeout_)); 129 : } 130 : 131 25 : regex XmppSession::tag_to_pattern(const char *tag) { 132 25 : std::string token("</"); 133 25 : token += ++tag; 134 25 : token += "[\\s\\t\\r\\n]*>"; 135 : 136 50 : regex exp(token.c_str()); 137 50 : return exp; 138 25 : } 139 : 140 117 : void XmppSession::SetBuf(const std::string &str) { 141 117 : if (buf_.empty()) { 142 104 : ReplaceBuf(str); 143 : } else { 144 13 : int pos = offset_ - buf_.begin(); 145 13 : buf_ += str; 146 13 : offset_ = buf_.begin() + pos; 147 : } 148 117 : } 149 : 150 138 : void XmppSession::ReplaceBuf(const std::string &str) { 151 138 : buf_ = str; 152 138 : buf_.reserve(kMaxMessageSize+8); 153 138 : offset_ = buf_.begin(); 154 138 : } 155 : 156 84 : bool XmppSession::LeftOver() const { 157 84 : if (buf_.empty()) 158 0 : return false; 159 84 : return (buf_.end() != offset_); 160 : } 161 : 162 : // Match a pattern in the buffer. Partially matched string is 163 : // kept in buf_ for use in conjucntion with next buffer read. 164 130 : int XmppSession::MatchRegex(const regex &patt) { 165 : 166 130 : std::string::const_iterator end = buf_.end(); 167 : 168 130 : if (regex_search(offset_, end, res_, patt, 169 130 : boost::match_default | boost::match_partial) == 0) { 170 3 : return -1; 171 : } 172 127 : if(res_[0].matched == false) { 173 : // partial match 174 11 : offset_ = res_[0].first; 175 11 : return 1; 176 : } else { 177 116 : begin_tag_ = string(res_[0].first, res_[0].second); 178 116 : offset_ = res_[0].second; 179 116 : return 0; 180 : } 181 : } 182 : 183 139 : bool XmppSession::Match(Buffer buffer, int *result, bool NewBuf) { 184 139 : const XmppConnection *connection = this->Connection(); 185 : 186 139 : if (connection == NULL) { 187 0 : return true; 188 : } 189 : 190 139 : xmsm::XmState state = connection->GetStateMcState(); 191 : xmsm::XmOpenConfirmState oc_state = 192 139 : connection->GetStateMcOpenConfirmState(); 193 : 194 139 : if (NewBuf) { 195 114 : const uint8_t *cp = BufferData(buffer); 196 : // TODO Avoid this copy 197 114 : std::string str(cp, cp + BufferSize(buffer)); 198 114 : XmppSession::SetBuf(str); 199 114 : } 200 : 201 139 : int m = -1; 202 139 : *result = 0; 203 : do { 204 193 : if (!tag_known_) { 205 : // check for whitespaces 206 138 : size_t pos = buf_.find_first_not_of(sXMPP_VALIDWS); 207 138 : if (pos != 0) { 208 75 : if (pos == string::npos) pos = buf_.size(); 209 75 : offset_ = buf_.begin() + pos; 210 75 : return false; 211 : } 212 : } 213 : 214 118 : if (state == xmsm::ACTIVE || state == xmsm::IDLE) { 215 26 : m = MatchRegex(tag_known_ ? stream_res_end_:stream_patt_); 216 92 : } else if (state == xmsm::CONNECT || state == xmsm::OPENSENT) { 217 : // Note, these are client only states 218 28 : if (!stream_open_matched_) { 219 26 : m = MatchRegex(tag_known_ ? stream_res_end_:stream_patt_); 220 26 : if ((m == 0) && (tag_known_)) { 221 13 : stream_open_matched_ = true; 222 : } 223 : } else { 224 2 : m = MatchRegex(tag_known_ ? tag_to_pattern(begin_tag_.c_str()): 225 : stream_features_patt_); 226 : } 227 64 : } else if ((state == xmsm::OPENCONFIRM) && !(IsSslDisabled())) { 228 8 : if (connection->IsClient()) { 229 4 : if (oc_state == xmsm::OPENCONFIRM_FEATURE_NEGOTIATION) { 230 2 : m = MatchRegex(tag_known_ ? end_patt_: proceed_patt_); 231 2 : if ((m == 0) && (tag_known_)) { 232 : // set the flag, as we do not want OnRead function to 233 : // read any more data from basic socket. 234 1 : SetSslHandShakeInProgress(true); 235 : } 236 2 : } else if (oc_state == xmsm::OPENCONFIRM_FEATURE_SUCCESS) { 237 2 : m = MatchRegex(tag_known_ ? stream_res_end_:stream_patt_); 238 : } else { 239 0 : m = MatchRegex(tag_known_ ? tag_to_pattern(begin_tag_.c_str()): 240 : stream_features_patt_); 241 : } 242 : } else { 243 4 : if (oc_state == xmsm::OPENCONFIRM_FEATURE_SUCCESS) { 244 2 : m = MatchRegex(tag_known_ ? stream_res_end_:stream_patt_); 245 : } else { 246 2 : m = MatchRegex(tag_known_ ? end_patt_:starttls_patt_); 247 2 : if ((m == 0) && (tag_known_)) { 248 1 : SetSslHandShakeInProgress(true); 249 : } 250 : } 251 : } 252 56 : } else if (state == xmsm::OPENCONFIRM || state == xmsm::ESTABLISHED) { 253 56 : m = MatchRegex(tag_known_ ? tag_to_pattern(begin_tag_.c_str()):patt_); 254 : } 255 : 256 118 : if (m == 0) { // full match 257 108 : *result = 0; 258 108 : tag_known_ ^= 1; 259 108 : if (!tag_known_) { 260 : // Found well formed xml 261 54 : return false; 262 : } 263 10 : } else if (m == -1) { // no match 264 1 : return true; 265 : } else { 266 9 : return true; // partial. read more 267 : } 268 54 : } while (true); 269 : 270 : return true; 271 : } 272 : 273 : // Read the socket stream and send messages to the connection object. 274 : // The buffer is copied to local string for regex match. 275 : // TODO Code need to change st Match() is done on buffer itself. 276 114 : void XmppSession::OnRead(Buffer buffer) { 277 114 : if (this->Connection() == NULL || !connection_) { 278 : // Connection is deleted. Session is being deleted as well 279 : // Drop the packet. 280 0 : ReleaseBuffer(buffer); 281 0 : return; 282 : } 283 : 284 114 : if (connection_->disable_read()) { 285 0 : ReleaseBuffer(buffer); 286 : 287 : // Reset the hold timer as we did receive some thing from the peer 288 0 : connection_->state_machine()->StartHoldTimer(); 289 0 : return; 290 : } 291 : 292 114 : int result = 0; 293 114 : bool more = Match(buffer, &result, true); 294 : do { 295 139 : if (more == false) { 296 129 : if (result < 0) { 297 : // TODO generate error, close connection. 298 45 : break; 299 : } 300 : 301 : // We got good match. Process the message 302 129 : std::string::const_iterator st = buf_.begin(); 303 129 : std::string xml = string(st, offset_); 304 : // Ensure we have not reached the end 305 129 : if (buf_.begin() == offset_) { // xml.size() == 0 306 45 : buf_.clear(); 307 45 : break; 308 : } 309 : 310 84 : connection_->ReceiveMsg(this, xml); 311 : 312 129 : } else { 313 : // Read more data. Either we have partial match 314 : // or no match but in this state we need to keep 315 : // reading data. 316 10 : break; 317 : } 318 : 319 84 : if (LeftOver()) { 320 25 : std::string::const_iterator st = buf_.end(); 321 25 : ReplaceBuf(string(offset_, st)); 322 25 : more = Match(buffer, &result, false); 323 : } else { 324 : // No more data in the Buffer 325 59 : buf_.clear(); 326 59 : break; 327 : } 328 25 : } while (true); 329 : 330 114 : ReleaseBuffer(buffer); 331 114 : return; 332 : }