- log(DEBUG,"Fetching a buffered packet size %d",buffer.size());
- packet_buf pb;
- pb.p.id = p.id;
- pb.p.key = p.key;
- pb.p.type = p.type;
- strcpy(pb.p.data,p.data);
- strcpy(pb.host,inet_ntoa(host_address.sin_addr));
- pb.port = ntohs(host_address.sin_port);
- this->buffer.push_back(pb);
-
- strcpy(message,buffer[0].p.data);
- theirkey = buffer[0].p.key;
- strcpy(host,buffer[0].host);
- prt = buffer[0].port;
-
- buffer.erase(buffer.begin());
+ // returns false if the packet could not be sent (e.g. target host down)
+ int rcvsize = 0;
+
+ // check if theres any data on this socket
+ // if not, continue onwards to the next.
+ pollfd polls;
+ polls.fd = this->connectors[i].GetDescriptor();
+ polls.events = POLLIN;
+ int ret = poll(&polls,1,1);
+ if (ret <= 0) continue;
+
+ rcvsize = recv(this->connectors[i].GetDescriptor(),data,65000,0);
+ data[rcvsize] = '\0';
+ if (rcvsize == -1)
+ {
+ if (errno != EAGAIN)
+ {
+ log(DEBUG,"recv() failed for Connection::RecvPacket(): %s",strerror(errno));
+ log(DEBUG,"Disabling connector: %s",this->connectors[i].GetServerName().c_str());
+ this->connectors[i].CloseConnection();
+ this->connectors[i].SetState(STATE_DISCONNECTED);
+ }
+ }
+ int pushed = 0;
+ if (rcvsize > 0)
+ {
+ this->connectors[i].AddBuffer(data);
+ if (this->connectors[i].BufferIsComplete())
+ {
+ while (this->connectors[i].BufferIsComplete())
+ {
+ std::string text = this->connectors[i].GetBuffer();
+ if (text != "")
+ {
+ if ((text[0] == ':') && (text.find(" ") != std::string::npos))
+ {
+ std::string orig = text;
+ log(DEBUG,"Original: %s",text.c_str());
+ std::string sum = text.substr(1,text.find(" ")-1);
+ text = text.substr(text.find(" ")+1,text.length());
+ std::string possible_token = text.substr(1,text.find(" ")-1);
+ if (possible_token.length() > 1)
+ {
+ sums.push_back("*");
+ text = orig;
+ log(DEBUG,"Non-mesh, non-tokenized string passed up the chain");
+ }
+ else
+ {
+ log(DEBUG,"Packet sum: '%s'",sum.c_str());
+ if ((already_have_sum(sum)) && (sum != "*"))
+ {
+ // we don't accept dupes
+ log(DEBUG,"Duplicate packet sum %s from server %s dropped",sum.c_str(),this->connectors[i].GetServerName().c_str());
+ continue;
+ }
+ sums.push_back(sum.c_str());
+ }
+ }
+ else sums.push_back("*");
+ messages.push_back(text.c_str());
+ strlcpy(recvhost,this->connectors[i].GetServerName().c_str(),160);
+ log(DEBUG,"Connection::RecvPacket() %d:%s->%s",pushed++,recvhost,text.c_str());
+ }
+ }
+ return true;
+ }
+ }