|
|
|
@ -570,42 +570,42 @@ bool CNode::ReceiveMsgBytes(const char *pch, unsigned int nBytes, bool& complete
|
|
|
|
|
nLastRecv = nTimeMicros / 1000000;
|
|
|
|
|
nRecvBytes += nBytes;
|
|
|
|
|
while (nBytes > 0) {
|
|
|
|
|
|
|
|
|
|
// get current incomplete message, or create a new one
|
|
|
|
|
if (vRecvMsg.empty() ||
|
|
|
|
|
vRecvMsg.back().complete())
|
|
|
|
|
vRecvMsg.push_back(CNetMessage(Params().MessageStart(), SER_NETWORK, INIT_PROTO_VERSION));
|
|
|
|
|
|
|
|
|
|
CNetMessage& msg = vRecvMsg.back();
|
|
|
|
|
|
|
|
|
|
// absorb network data
|
|
|
|
|
int handled;
|
|
|
|
|
if (!msg.in_data)
|
|
|
|
|
handled = msg.readHeader(pch, nBytes);
|
|
|
|
|
if (!m_deserializer->in_data)
|
|
|
|
|
handled = m_deserializer->readHeader(pch, nBytes);
|
|
|
|
|
else
|
|
|
|
|
handled = msg.readData(pch, nBytes);
|
|
|
|
|
handled = m_deserializer->readData(pch, nBytes);
|
|
|
|
|
|
|
|
|
|
if (handled < 0)
|
|
|
|
|
if (handled < 0) {
|
|
|
|
|
m_deserializer->Reset();
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (msg.in_data && msg.hdr.nMessageSize > MAX_PROTOCOL_MESSAGE_LENGTH) {
|
|
|
|
|
if (m_deserializer->in_data && m_deserializer->hdr.nMessageSize > MAX_PROTOCOL_MESSAGE_LENGTH) {
|
|
|
|
|
LogPrint(BCLog::NET, "Oversized message from peer=%i, disconnecting\n", GetId());
|
|
|
|
|
m_deserializer->Reset();
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pch += handled;
|
|
|
|
|
nBytes -= handled;
|
|
|
|
|
|
|
|
|
|
if (msg.complete()) {
|
|
|
|
|
if (m_deserializer->complete()) {
|
|
|
|
|
// decompose a transport agnostic CNetMessage from the deserializer
|
|
|
|
|
CNetMessage msg = m_deserializer->GetMessage(Params().MessageStart(), nTimeMicros);
|
|
|
|
|
|
|
|
|
|
//store received bytes per message command
|
|
|
|
|
//to prevent a memory DOS, only allow valid commands
|
|
|
|
|
mapMsgCmdSize::iterator i = mapRecvBytesPerMsgCmd.find(msg.hdr.pchCommand);
|
|
|
|
|
mapMsgCmdSize::iterator i = mapRecvBytesPerMsgCmd.find(m_deserializer->hdr.pchCommand);
|
|
|
|
|
if (i == mapRecvBytesPerMsgCmd.end())
|
|
|
|
|
i = mapRecvBytesPerMsgCmd.find(NET_MESSAGE_COMMAND_OTHER);
|
|
|
|
|
assert(i != mapRecvBytesPerMsgCmd.end());
|
|
|
|
|
i->second += msg.hdr.nMessageSize + CMessageHeader::HEADER_SIZE;
|
|
|
|
|
i->second += m_deserializer->hdr.nMessageSize + CMessageHeader::HEADER_SIZE;
|
|
|
|
|
|
|
|
|
|
// push the message to the process queue,
|
|
|
|
|
vRecvMsg.push_back(std::move(msg));
|
|
|
|
|
|
|
|
|
|
msg.nTime = nTimeMicros;
|
|
|
|
|
complete = true;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@ -639,8 +639,7 @@ int CNode::GetSendVersion() const
|
|
|
|
|
return nSendVersion;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
int CNetMessage::readHeader(const char *pch, unsigned int nBytes)
|
|
|
|
|
int TransportDeserializer::readHeader(const char *pch, unsigned int nBytes)
|
|
|
|
|
{
|
|
|
|
|
// copy data to temporary parsing buffer
|
|
|
|
|
unsigned int nRemaining = 24 - nHdrPos;
|
|
|
|
@ -671,7 +670,7 @@ int CNetMessage::readHeader(const char *pch, unsigned int nBytes)
|
|
|
|
|
return nCopy;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int CNetMessage::readData(const char *pch, unsigned int nBytes)
|
|
|
|
|
int TransportDeserializer::readData(const char *pch, unsigned int nBytes)
|
|
|
|
|
{
|
|
|
|
|
unsigned int nRemaining = hdr.nMessageSize - nDataPos;
|
|
|
|
|
unsigned int nCopy = std::min(nRemaining, nBytes);
|
|
|
|
@ -688,7 +687,7 @@ int CNetMessage::readData(const char *pch, unsigned int nBytes)
|
|
|
|
|
return nCopy;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const uint256& CNetMessage::GetMessageHash() const
|
|
|
|
|
const uint256& TransportDeserializer::GetMessageHash() const
|
|
|
|
|
{
|
|
|
|
|
assert(complete());
|
|
|
|
|
if (data_hash.IsNull())
|
|
|
|
@ -696,6 +695,35 @@ const uint256& CNetMessage::GetMessageHash() const
|
|
|
|
|
return data_hash;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
CNetMessage TransportDeserializer::GetMessage(const CMessageHeader::MessageStartChars& message_start, int64_t time) {
|
|
|
|
|
// decompose a single CNetMessage from the TransportDeserializer
|
|
|
|
|
CNetMessage msg(std::move(vRecv));
|
|
|
|
|
|
|
|
|
|
// store state about valid header, netmagic and checksum
|
|
|
|
|
msg.m_valid_header = hdr.IsValid(message_start);
|
|
|
|
|
msg.m_valid_netmagic = (memcmp(hdr.pchMessageStart, message_start, CMessageHeader::MESSAGE_START_SIZE) == 0);
|
|
|
|
|
uint256 hash = GetMessageHash();
|
|
|
|
|
|
|
|
|
|
// store command string, payload size
|
|
|
|
|
msg.m_command = hdr.GetCommand();
|
|
|
|
|
msg.m_message_size = hdr.nMessageSize;
|
|
|
|
|
|
|
|
|
|
msg.m_valid_checksum = (memcmp(hash.begin(), hdr.pchChecksum, CMessageHeader::CHECKSUM_SIZE) == 0);
|
|
|
|
|
if (!msg.m_valid_checksum) {
|
|
|
|
|
LogPrint(BCLog::NET, "CHECKSUM ERROR (%s, %u bytes), expected %s was %s\n",
|
|
|
|
|
SanitizeString(msg.m_command), msg.m_message_size,
|
|
|
|
|
HexStr(hash.begin(), hash.begin()+CMessageHeader::CHECKSUM_SIZE),
|
|
|
|
|
HexStr(hdr.pchChecksum, hdr.pchChecksum+CMessageHeader::CHECKSUM_SIZE));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// store receive time
|
|
|
|
|
msg.m_time = time;
|
|
|
|
|
|
|
|
|
|
// reset the network deserializer (prepare for the next message)
|
|
|
|
|
Reset();
|
|
|
|
|
return msg;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
size_t CConnman::SocketSendData(CNode *pnode) const EXCLUSIVE_LOCKS_REQUIRED(pnode->cs_vSend)
|
|
|
|
|
{
|
|
|
|
|
auto it = pnode->vSendMsg.begin();
|
|
|
|
@ -1347,9 +1375,9 @@ void CConnman::SocketHandler()
|
|
|
|
|
size_t nSizeAdded = 0;
|
|
|
|
|
auto it(pnode->vRecvMsg.begin());
|
|
|
|
|
for (; it != pnode->vRecvMsg.end(); ++it) {
|
|
|
|
|
if (!it->complete())
|
|
|
|
|
break;
|
|
|
|
|
nSizeAdded += it->vRecv.size() + CMessageHeader::HEADER_SIZE;
|
|
|
|
|
// vRecvMsg contains only completed CNetMessage
|
|
|
|
|
// the single possible partially deserialized message are held by TransportDeserializer
|
|
|
|
|
nSizeAdded += it->m_recv.size() + CMessageHeader::HEADER_SIZE;
|
|
|
|
|
}
|
|
|
|
|
{
|
|
|
|
|
LOCK(pnode->cs_vProcessMsg);
|
|
|
|
@ -2678,6 +2706,8 @@ CNode::CNode(NodeId idIn, ServiceFlags nLocalServicesIn, int nMyStartingHeightIn
|
|
|
|
|
} else {
|
|
|
|
|
LogPrint(BCLog::NET, "Added connection peer=%d\n", id);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
m_deserializer = MakeUnique<TransportDeserializer>(TransportDeserializer(Params().MessageStart(), SER_NETWORK, INIT_PROTO_VERSION));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
CNode::~CNode()
|
|
|
|
|