diff --git a/src/friend_server/fsclient.cc b/src/friend_server/fsclient.cc index 85eb7ffb8..c25808e1a 100644 --- a/src/friend_server/fsclient.cc +++ b/src/friend_server/fsclient.cc @@ -202,15 +202,17 @@ bool FsClient::sendItem(const std::string& server_address,uint16_t server_port, uint32_t ss; p.SendItem(item,ss); time_t now = time(nullptr); + bool got_response = false; - while(now > time(nullptr)+15) // wait 15 secs for a response. +#ifdef DEBUG_FSCLIENT + RsDbg() << "Ticking for response..."; +#endif + while(now + 15 > time(nullptr)) // wait 15 secs for a response. { p.tick(); // ticks bio RsItem *ritem = GetItem(); -#ifdef DEBUG_FSCLIENT - RsDbg() << "Ticking for response..."; -#endif + if(ritem) { response.push_back(ritem); @@ -218,12 +220,16 @@ bool FsClient::sendItem(const std::string& server_address,uint16_t server_port, std::cerr << *ritem << std::endl; RsDbg() << "End of transmission. " ; + got_response = true; break; } else std::this_thread::sleep_for(std::chrono::milliseconds(200)); } + if(!got_response) + std::cerr << "Sending timed out. Connection is dead?" << std::endl; + RsDbg() << " Stopping/killing pqistreamer" ; p.fullstop(); diff --git a/src/pqi/pqi_base.h b/src/pqi/pqi_base.h index 72d5d540b..e0ca235a0 100644 --- a/src/pqi/pqi_base.h +++ b/src/pqi/pqi_base.h @@ -286,6 +286,11 @@ public: /** * reads data from a prescribed location (implementation dependent) + * -- WARNING -- if used to feed a pqistreamer, the streamer will assume one of the two situations: + * 1 - len bytes are read, or + * 2 - less than len bytes are read, but in this case the interface should keep the data + * This rule corresponds to openssl behavior. If not followed, incoming packets will be lost. + * *@param data what will be sent *@param len the size of data pointed to in memory */ diff --git a/src/pqi/pqifdbin.cc b/src/pqi/pqifdbin.cc index 1d3fd41ce..b3e81d8cd 100644 --- a/src/pqi/pqifdbin.cc +++ b/src/pqi/pqifdbin.cc @@ -24,6 +24,8 @@ #include "util/rsfile.h" #include "pqi/pqifdbin.h" +#define DEBUG_FS_BIN + #ifdef DEBUG_FS_BIN #include "util/rsprint.h" #endif @@ -103,7 +105,7 @@ int RsFdBinInterface::read_pending() if(readbytes == 0) { RsDbg() << "Reached END of the stream!" ; - RsDbg() << "Closing socket!" ; + RsDbg() << "Closing socket! mTotalInBufferBytes = " << mTotalInBufferBytes ; close(); return mTotalInBufferBytes; @@ -129,8 +131,8 @@ int RsFdBinInterface::read_pending() if(readbytes > 0) { #ifdef DEBUG_FS_BIN - RsDbg() << "Received the following bytes: size=" << readbytes << " len=" << RsUtil::BinToHex( reinterpret_cast(inBuffer),readbytes,50) << std::endl; - RsDbg() << "Received the following bytes: size=" << readbytes << " len=" << std::string(inBuffer,readbytes) << std::endl; + RsDbg() << "Received the following bytes: size=" << readbytes << " data=" << RsUtil::BinToHex( reinterpret_cast(inBuffer),readbytes,50) << std::endl; + //RsDbg() << "Received the following bytes: size=" << readbytes << " len=" << std::string(inBuffer,readbytes) << std::endl; #endif void *ptr = malloc(readbytes); @@ -148,6 +150,7 @@ int RsFdBinInterface::read_pending() RsDbg() << "Socket: " << mCLintConnt << ". Total read: " << mTotalReadBytes << ". Buffer size: " << mTotalInBufferBytes ; #endif } + RsDbg() << "End of read_pending: mTotalInBufferBytes = " << mTotalInBufferBytes; return mTotalInBufferBytes; } @@ -164,6 +167,8 @@ int RsFdBinInterface::write_pending() written = send(mCLintConnt, (char*) p.first, p.second, 0); else #endif + std::cerr << "RsFdBinInterface -- SENDING --- len=" << p.second << " data=" << RsUtil::BinToHex((uint8_t*)p.first,p.second)<< std::endl; + written = write(mCLintConnt, p.first, p.second); if(written < 0) @@ -217,6 +222,7 @@ int RsFdBinInterface::write_pending() RsFdBinInterface::~RsFdBinInterface() { + close(); clean(); } @@ -243,37 +249,44 @@ int RsFdBinInterface::readline(void *data, int len) int RsFdBinInterface::readdata(void *data, int len) { + // Expected behavior of BinInterface: when the full amount of bytes (len bytes) can not be provided, we keep the data. + + if((int)mTotalInBufferBytes < len) + { + std::cerr << "RsFdBinInterface -- READ --- not enough data to fill " << len << " bytes. Current buffer is " << mTotalInBufferBytes << " bytes." << std::endl; + return mTotalInBufferBytes; + } + // read incoming bytes in the buffer int total_len = 0; while(total_len < len) { - if(in_buffer.empty()) - { - mTotalInBufferBytes -= total_len; -#ifdef DEBUG_FS_BIN - std::cerr << "RsFdBinInterface -- READ --- len=" << total_len << " data=" << RsUtil::BinToHex((uint8_t*)data,total_len)<< std::endl; -#endif - return total_len; - } - // If the remaining buffer is too large, chop of the beginning of it. - if(total_len + in_buffer.front().second > len) + if(total_len + in_buffer.front().second >= len) { int bytes_in = len - total_len; int bytes_out = in_buffer.front().second - bytes_in; memcpy(&(static_cast(data)[total_len]),in_buffer.front().first,bytes_in); - void *ptr = malloc(bytes_out); - memcpy(ptr,&(static_cast(in_buffer.front().first)[bytes_in]),bytes_out); + auto tmp_ptr = in_buffer.front().first; - free(in_buffer.front().first); - in_buffer.front().first = ptr; + if(bytes_out > 0) + { + void *ptr = malloc(bytes_out); + memcpy(ptr,&(static_cast(in_buffer.front().first)[bytes_in]),bytes_out); + in_buffer.front().first = ptr; + } + + free(tmp_ptr); in_buffer.front().second -= bytes_in; + if(in_buffer.front().second == 0) + in_buffer.pop_front(); + mTotalInBufferBytes -= len; #ifdef DEBUG_FS_BIN std::cerr << "RsFdBinInterface -- READ --- len=" << len << " data=" << RsUtil::BinToHex((uint8_t*)data,len)<< std::endl; @@ -290,11 +303,10 @@ int RsFdBinInterface::readdata(void *data, int len) in_buffer.pop_front(); } } - mTotalInBufferBytes -= len; #ifdef DEBUG_FS_BIN - std::cerr << "RsFdBinInterface -- READ --- len=" << len << " data=" << RsUtil::BinToHex((uint8_t*)data,len)<< std::endl; + std::cerr << "RsFdBinInterface: ERROR. Shouldn't be here!" << std::endl; #endif - return len; + return 0; } int RsFdBinInterface::senddata(void *data, int len) @@ -302,7 +314,7 @@ int RsFdBinInterface::senddata(void *data, int len) // shouldn't we better send in multiple packets, similarly to how we read? #ifdef DEBUG_FS_BIN - std::cerr << "RsFdBinInterface -- SENDING --- len=" << len << " data=" << RsUtil::BinToHex((uint8_t*)data,len)<< std::endl; + std::cerr << "RsFdBinInterface -- QUEUEING OUT --- len=" << len << " data=" << RsUtil::BinToHex((uint8_t*)data,len)<< std::endl; #endif if(len == 0) { @@ -330,7 +342,7 @@ int RsFdBinInterface::netstatus() int RsFdBinInterface::isactive() { - return mIsActive ; + return mIsActive || mTotalInBufferBytes>0; } bool RsFdBinInterface::moretoread(uint32_t /* usec */) diff --git a/src/pqi/pqistreamer.cc b/src/pqi/pqistreamer.cc index 7f5a42301..516f6d9f3 100644 --- a/src/pqi/pqistreamer.cc +++ b/src/pqi/pqistreamer.cc @@ -40,6 +40,8 @@ #include "util/rsprint.h" // for BinToHex #include "util/rsstring.h" // for rs_sprintf_append, rs_sprintf +#define DEBUG_PQISTREAMER 1 + static struct RsLog::logInfo pqistreamerzoneInfo = {RsLog::Default, "pqistreamer"}; #define pqistreamerzone &pqistreamerzoneInfo @@ -898,8 +900,8 @@ continue_packet: int tmplen ; // Don't reset the block now! If pqissl is in the middle of a multiple-chunk - // packet (larger than 16384 bytes), and pqistreamer jumped directly yo - // continue_packet:, then readdata is going to write after the beginning of + // packet (larger than 16384 bytes), and pqistreamer jumped directly to + // continue_packet, then readdata is going to write after the beginning of // extradata, yet not exactly at start -> the start of the packet would be wiped out. // // so, don't do that: