Line data Source code
1 : /* 2 : * Each object of NexthopDBClient class interacts with the master 3 : * instance of NexthopDBServer and to the underlying (Unix socket) session 4 : * connected to the client process: 5 : * 6 : * {RemoveNexthop()} 7 : * {AddNexthop()} {Send()} 8 : * [NexthopDBServer] ------ [NexthopDBClient] ----- [UnixDomainSocketSession] 9 : * 1:n 1:1 10 : */ 11 : #include <base/logging.h> 12 : #include "nexthop_server/nexthop_server.h" 13 : #include <pthread.h> 14 : #include "rapidjson/document.h" 15 : #include "rapidjson/stringbuffer.h" 16 : #include "rapidjson/writer.h" 17 : 18 0 : NexthopDBClient::NexthopDBClient(UnixDomainSocketSession *session, 19 0 : NexthopDBServer *server) 20 0 : : nexthop_list_() 21 : { 22 0 : server_ = server; 23 0 : session_ = session; 24 0 : session_->set_observer(boost::bind(&NexthopDBClient::EventHandler, this, 25 : _1, _2)); 26 0 : } 27 : 28 0 : NexthopDBClient::~NexthopDBClient() 29 : { 30 0 : for (NexthopListIterator iter = nexthop_list_.begin(); 31 0 : iter != nexthop_list_.end(); ++iter) { 32 0 : nexthop_list_.erase (iter); 33 : } 34 0 : } 35 : 36 : void 37 0 : NexthopDBClient::EventHandler(UnixDomainSocketSession *session, 38 : UnixDomainSocketSession::Event event) 39 : { 40 0 : if (event == UnixDomainSocketSession::READY) { 41 0 : WriteMessage(); 42 0 : } else if (event == UnixDomainSocketSession::CLOSE) { 43 0 : LOG (DEBUG, "[NexthopServer] Client " << session->session_id() << 44 : " closed session"); 45 0 : server_->RemoveClient(session->session_id()); 46 : } 47 0 : } 48 : 49 : void 50 0 : NexthopDBClient::AddNexthop(NexthopDBEntry::NexthopPtr nh) 51 : { 52 0 : nexthop_list_.push_back(nh); 53 0 : } 54 : 55 : void 56 0 : NexthopDBClient::RemoveNexthop(NexthopDBEntry::NexthopPtr nh) 57 : { 58 0 : for (NexthopListIterator iter = nexthop_list_.begin(); 59 0 : iter != nexthop_list_.end(); ++iter) { 60 0 : if (*(*iter) == *nh) { 61 0 : iter = nexthop_list_.erase(iter); 62 0 : break; 63 : } 64 : } 65 0 : } 66 : 67 : bool 68 0 : NexthopDBClient::FindNexthop(NexthopDBEntry::NexthopPtr nh) 69 : { 70 0 : for (NexthopListIterator iter = nexthop_list_.begin(); 71 0 : iter != nexthop_list_.end(); ++iter) { 72 0 : if (*(*iter) == *nh) { 73 0 : return true; 74 : } 75 : } 76 0 : return false; 77 : } 78 : 79 : /* Caller should free the returned buffer */ 80 : uint8_t * 81 0 : NexthopDBClient::NextMessage(int *data_len) 82 : { 83 0 : contrail_rapidjson::StringBuffer s; 84 0 : contrail_rapidjson::Writer <contrail_rapidjson::StringBuffer> writer(s); 85 : 86 0 : writer.StartObject(); 87 : 88 0 : int pdu_len = 0; 89 : 90 0 : for (NexthopListIterator iter = nexthop_list_.begin(); 91 0 : iter != nexthop_list_.end();) { 92 0 : NexthopDBEntry::NexthopPtr tnh = *iter; 93 : 94 0 : if ((pdu_len + tnh->EncodedLength()) > 95 : UnixDomainSocketSession::kPDUDataLen) { 96 0 : break; 97 : } 98 : 99 0 : writer.String((tnh->nexthop_string()).c_str()); 100 0 : writer.StartObject(); 101 0 : const char *action = (tnh->state() == 102 0 : NexthopDBEntry::NEXTHOP_STATE_DELETED) ? "del" : "add"; 103 0 : writer.String("action"); 104 0 : writer.String(action); 105 0 : writer.EndObject(); 106 : 107 0 : pdu_len += tnh->EncodedLength(); 108 : 109 0 : iter = nexthop_list_.erase(iter); 110 0 : } 111 : 112 0 : writer.EndObject(); 113 : 114 : /* Nothing to consume? */ 115 0 : if (!pdu_len) { 116 0 : return NULL; 117 : } 118 : 119 0 : const char *nhdata = s.GetString(); 120 0 : int nhlen = strlen(nhdata); 121 0 : u_int8_t *out_data = new u_int8_t[nhlen + 4]; 122 0 : out_data[0] = (unsigned char) (nhlen >> 24); 123 0 : out_data[1] = (unsigned char) (nhlen >> 16); 124 0 : out_data[2] = (unsigned char) (nhlen >> 8); 125 0 : out_data[3] = (unsigned char) nhlen; 126 0 : memcpy(&out_data[4], nhdata, nhlen); 127 0 : *data_len = nhlen + 4; 128 0 : return out_data; 129 0 : } 130 : 131 : void 132 0 : NexthopDBClient::WriteMessage() 133 : { 134 0 : uint8_t *data = NULL; 135 : int data_len; 136 : 137 : while(1) { 138 0 : data_len = 0; 139 0 : data = NextMessage(&data_len); 140 0 : if (!data || !data_len) { 141 : break; 142 : } 143 : 144 : /* 145 : * Send always succeeds 146 : */ 147 0 : session_->Send(data, data_len); 148 0 : delete [] data; 149 : } 150 0 : }