19 #include "visiontransfer/datablockprotocol.h" 20 #include "visiontransfer/exceptions.h" 26 #include <arpa/inet.h> 29 #define LOG_ERROR(expr) 35 : isServer(server), protType(protType),
36 transferDone(true), rawData(nullptr), rawValidBytes(0),
37 transferOffset(0), transferSize(0), overwrittenTransferData(0),
38 overwrittenTransferIndex(-1), transferHeaderData(nullptr),
39 transferHeaderSize(0), waitingForMissingSegments(false),
40 totalReceiveSize(0), connectionConfirmed(false),
41 confirmationMessagePending(false), eofMessagePending(false),
42 clientConnectionPending(false), resendMessagePending(false),
43 lastRemoteHostActivity(), lastSentHeartbeat(),
44 lastReceivedHeartbeat(
std::chrono::steady_clock::now()),
45 receiveOffset(0), finishedReception(false), droppedReceptions(0),
46 unprocessedMsgLength(0), headerReceived(false) {
48 if(protType == PROTOCOL_TCP) {
49 maxPayloadSize = MAX_TCP_BYTES_TRANSFER;
52 maxPayloadSize = maxUdpPacketSize -
sizeof(int);
53 minPayloadSize = maxPayloadSize;
55 resizeReceiveBuffer();
60 overwrittenTransferIndex = -1;
63 missingTransferSegments.clear();
67 if(!transferDone && transferOffset > 0) {
69 }
else if(headerSize + 9 > static_cast<int>(
sizeof(controlMessageBuffer))) {
74 this->transferSize = transferSize;
76 transferHeaderData = &data[-6];
78 unsigned short netHeaderSize = htons(static_cast<unsigned short>(headerSize));
79 memcpy(transferHeaderData, &netHeaderSize,
sizeof(netHeaderSize));
81 unsigned int netTransferSize = htonl(static_cast<unsigned int>(transferSize));
82 memcpy(&transferHeaderData[2], &netTransferSize,
sizeof(netTransferSize));
85 if(protType == PROTOCOL_UDP) {
87 transferHeaderData[headerSize++] = HEADER_MESSAGE;
88 transferHeaderData[headerSize++] = 0xFF;
89 transferHeaderData[headerSize++] = 0xFF;
90 transferHeaderData[headerSize++] = 0xFF;
91 transferHeaderData[headerSize++] = 0xFF;
94 transferHeaderSize = headerSize;
98 if(transferHeaderSize == 0 || transferHeaderData ==
nullptr) {
102 transferDone =
false;
105 overwrittenTransferIndex = -1;
106 rawValidBytes = min(transferSize, validBytes);
110 if(validBytes >= transferSize) {
111 rawValidBytes = transferSize;
112 }
else if(validBytes < static_cast<int>(
sizeof(
int))) {
115 rawValidBytes = validBytes;
120 if(transferDone || rawValidBytes == 0) {
127 if(protType == PROTOCOL_TCP && transferOffset == 0 && transferHeaderData !=
nullptr) {
128 length = transferHeaderSize;
129 const unsigned char* ret = transferHeaderData;
130 transferHeaderData =
nullptr;
136 restoreTransferBuffer();
140 getNextTransferSegment(offset, length);
145 if(protType == PROTOCOL_UDP) {
147 overwrittenTransferIndex = offset + length;
148 int* offsetPtr =
reinterpret_cast<int*
>(&rawData[offset + length]);
149 overwrittenTransferData = *offsetPtr;
150 *offsetPtr =
static_cast<int>(htonl(offset));
151 length +=
sizeof(int);
154 return &rawData[offset];
157 void DataBlockProtocol::getNextTransferSegment(
int& offset,
int& length) {
158 if(missingTransferSegments.size() == 0) {
160 length = min(maxPayloadSize, rawValidBytes - transferOffset);
161 if(length == 0 || (length < minPayloadSize && rawValidBytes != transferSize)) {
166 offset = transferOffset;
167 transferOffset += length;
169 if(transferOffset >= transferSize && protType == PROTOCOL_UDP) {
170 eofMessagePending =
true;
174 length = min(maxPayloadSize, missingTransferSegments.front().second);
175 offset = missingTransferSegments.front().first;
176 LOG_ERROR(
"Re-transmitting: " << offset <<
" - " << (offset + length));
178 int remaining = missingTransferSegments[0].second - length;
181 missingTransferSegments.pop_front();
184 missingTransferSegments.front().first += length;
185 missingTransferSegments.front().second = remaining;
190 void DataBlockProtocol::restoreTransferBuffer() {
191 if(overwrittenTransferIndex > 0) {
192 *
reinterpret_cast<int*
>(&rawData[overwrittenTransferIndex]) = overwrittenTransferData;
194 overwrittenTransferIndex = -1;
198 return transferOffset >= transferSize && !eofMessagePending;
202 if(protType == PROTOCOL_TCP) {
203 return MAX_TCP_BYTES_TRANSFER;
205 return MAX_UDP_RECEPTION;
210 if(static_cast<int>(receiveBuffer.size() - receiveOffset) < maxLength) {
214 return &receiveBuffer[receiveOffset];
218 transferComplete =
false;
223 if(finishedReception) {
228 if(protType == PROTOCOL_UDP) {
229 processReceivedUdpMessage(length, transferComplete);
231 processReceivedTcpMessage(length, transferComplete);
234 transferComplete = finishedReception;
237 void DataBlockProtocol::processReceivedUdpMessage(
int length,
bool&
transferComplete) {
238 if(length < static_cast<int>(
sizeof(
int)) ||
239 receiveOffset + length > static_cast<int>(receiveBuffer.size())) {
244 int segmentOffset = ntohl(*reinterpret_cast<int*>(
245 &receiveBuffer[receiveOffset + length -
sizeof(
int)]));
247 if(segmentOffset == static_cast<int>(0xFFFFFFFF)) {
249 processControlMessage(length);
250 }
else if(segmentOffset < 0) {
254 int payloadLength = length -
sizeof(int);
256 if(segmentOffset != receiveOffset) {
259 if(!waitingForMissingSegments && receiveOffset > 0 && segmentOffset > receiveOffset) {
261 LOG_ERROR(
"Missing segment: " << receiveOffset <<
" - " << segmentOffset
262 <<
" (" << missingReceiveSegments.size() <<
")");
264 MissingReceiveSegment missingSeg;
265 missingSeg.offset = receiveOffset;
266 missingSeg.length = segmentOffset - receiveOffset;
267 missingSeg.isEof =
false;
268 memcpy(missingSeg.subsequentData, &receiveBuffer[receiveOffset],
269 sizeof(missingSeg.subsequentData));
270 missingReceiveSegments.push_back(missingSeg);
273 memcpy(&receiveBuffer[segmentOffset], &receiveBuffer[receiveOffset], payloadLength);
274 receiveOffset = segmentOffset;
280 if(segmentOffset > 0 ) {
281 if(receiveOffset > 0) {
282 LOG_ERROR(
"Resend failed!");
286 LOG_ERROR(
"Missed EOF message!");
291 if(segmentOffset == 0) {
293 lastRemoteHostActivity = std::chrono::steady_clock::now();
297 receiveOffset = getNextUdpReceiveOffset(segmentOffset, payloadLength);
301 int DataBlockProtocol::getNextUdpReceiveOffset(
int lastSegmentOffset,
int lastSegmentSize) {
302 if(!waitingForMissingSegments) {
304 return lastSegmentOffset + lastSegmentSize;
307 MissingReceiveSegment& firstSeg = missingReceiveSegments.front();
308 if(lastSegmentOffset != firstSeg.offset) {
309 LOG_ERROR(
"Received invalid resend: " << lastSegmentOffset);
313 firstSeg.offset += lastSegmentSize;
314 firstSeg.length -= lastSegmentSize;
315 if(firstSeg.length == 0) {
316 if(!firstSeg.isEof) {
317 memcpy(&receiveBuffer[firstSeg.offset + firstSeg.length],
318 firstSeg.subsequentData,
sizeof(firstSeg.subsequentData));
320 missingReceiveSegments.pop_front();
323 if(missingReceiveSegments.size() == 0) {
324 waitingForMissingSegments =
false;
325 finishedReception =
true;
326 return min(totalReceiveSize, static_cast<int>(receiveBuffer.size()));
328 return missingReceiveSegments.front().offset;
334 void DataBlockProtocol::processReceivedTcpMessage(
int length,
bool&
transferComplete) {
337 if(unprocessedMsgLength != 0) {
338 if(length + unprocessedMsgLength > MAX_OUTSTANDING_BYTES) {
342 ::memmove(&receiveBuffer[unprocessedMsgLength], &receiveBuffer[0], length);
343 ::memcpy(&receiveBuffer[0], &unprocessedMsgPart[0], unprocessedMsgLength);
344 length += unprocessedMsgLength;
345 unprocessedMsgLength = 0;
349 if(!headerReceived) {
350 int totalHeaderSize = parseReceivedHeader(length, receiveOffset);
351 if(totalHeaderSize == 0) {
353 ::memcpy(unprocessedMsgPart, &receiveBuffer[0], length);
354 unprocessedMsgLength = length;
359 length -= totalHeaderSize;
363 ::memmove(&receiveBuffer[0], &receiveBuffer[totalHeaderSize], length);
369 if(receiveOffset + length > totalReceiveSize) {
370 int newLength =
static_cast<int>(totalReceiveSize - receiveOffset);
372 if(unprocessedMsgLength != 0 || length - newLength > MAX_OUTSTANDING_BYTES) {
376 unprocessedMsgLength = length - newLength;
377 ::memcpy(unprocessedMsgPart, &receiveBuffer[receiveOffset + newLength], unprocessedMsgLength);
383 receiveOffset += length;
385 if(receiveOffset == totalReceiveSize) {
387 finishedReception =
true;
391 int DataBlockProtocol::parseReceivedHeader(
int length,
int offset) {
392 constexpr
int headerExtraBytes = 6;
394 if(length < headerExtraBytes) {
398 unsigned short headerSize = ntohs(*reinterpret_cast<unsigned short*>(&receiveBuffer[offset]));
399 totalReceiveSize =
static_cast<int>(ntohl(*reinterpret_cast<unsigned int*>(&receiveBuffer[offset + 2])));
401 if(headerSize + headerExtraBytes > static_cast<int>(receiveBuffer.size())
402 || totalReceiveSize < 0 || headerSize + headerExtraBytes > length ) {
406 headerReceived =
true;
407 receivedHeader.assign(receiveBuffer.begin() + offset + headerExtraBytes,
408 receiveBuffer.begin() + offset + headerSize + headerExtraBytes);
409 resizeReceiveBuffer();
411 return headerSize + headerExtraBytes;
415 headerReceived =
false;
417 missingReceiveSegments.clear();
418 receivedHeader.clear();
419 waitingForMissingSegments =
false;
420 totalReceiveSize = 0;
421 finishedReception =
false;
428 length = receiveOffset;
429 if(missingReceiveSegments.size() > 0) {
430 length = min(length, missingReceiveSegments[0].offset);
432 return &receiveBuffer[0];
436 if(receivedHeader.size() > 0) {
437 length =
static_cast<int>(receivedHeader.size());
438 return &receivedHeader[0];
444 bool DataBlockProtocol::processControlMessage(
int length) {
445 if(length < static_cast<int>(
sizeof(
int) + 1)) {
449 int payloadLength = length -
sizeof(int) - 1;
451 switch(receiveBuffer[receiveOffset + payloadLength]) {
452 case CONFIRM_MESSAGE:
454 connectionConfirmed =
true;
456 case CONNECTION_MESSAGE:
458 connectionConfirmed =
true;
459 confirmationMessagePending =
true;
460 clientConnectionPending =
true;
463 lastReceivedHeartbeat = std::chrono::steady_clock::now();
465 case HEADER_MESSAGE: {
466 int offset = receiveOffset;
467 if(receiveOffset != 0) {
468 if(receiveOffset == totalReceiveSize) {
469 LOG_ERROR(
"No EOF message received!");
471 LOG_ERROR(
"Received header too late/early!");
475 if(parseReceivedHeader(payloadLength, offset) == 0) {
482 if(receiveOffset != 0) {
483 parseEofMessage(length);
486 case RESEND_MESSAGE: {
488 parseResendMessage(payloadLength);
491 case HEARTBEAT_MESSAGE:
493 lastReceivedHeartbeat = std::chrono::steady_clock::now();
504 if(protType == PROTOCOL_TCP) {
507 }
else if(connectionConfirmed) {
508 return !isServer || std::chrono::duration_cast<std::chrono::milliseconds>(
509 std::chrono::steady_clock::now() - lastReceivedHeartbeat).count()
510 < 2*HEARTBEAT_INTERVAL_MS;
517 if(protType == PROTOCOL_TCP) {
522 if(confirmationMessagePending) {
524 confirmationMessagePending =
false;
525 controlMessageBuffer[0] = CONFIRM_MESSAGE;
527 }
else if(!isServer && std::chrono::duration_cast<std::chrono::milliseconds>(
528 std::chrono::steady_clock::now() - lastRemoteHostActivity).count() > RECONNECT_TIMEOUT_MS) {
530 controlMessageBuffer[0] = CONNECTION_MESSAGE;
534 lastRemoteHostActivity = lastSentHeartbeat = std::chrono::steady_clock::now();
535 }
else if(transferHeaderData !=
nullptr &&
isConnected()) {
537 length = transferHeaderSize;
538 const unsigned char* ret = transferHeaderData;
539 transferHeaderData =
nullptr;
541 }
else if(eofMessagePending) {
543 eofMessagePending =
false;
544 unsigned int networkOffset = htonl(static_cast<unsigned int>(transferOffset));
545 memcpy(&controlMessageBuffer[0], &networkOffset,
sizeof(
int));
546 controlMessageBuffer[
sizeof(int)] = EOF_MESSAGE;
548 }
else if(resendMessagePending) {
550 resendMessagePending =
false;
551 if(!generateResendRequest(length)) {
555 }
else if(!isServer && std::chrono::duration_cast<std::chrono::milliseconds>(
556 std::chrono::steady_clock::now() - lastSentHeartbeat).count() > HEARTBEAT_INTERVAL_MS) {
558 controlMessageBuffer[0] = HEARTBEAT_MESSAGE;
560 lastSentHeartbeat = std::chrono::steady_clock::now();
566 controlMessageBuffer[length++] = 0xff;
567 controlMessageBuffer[length++] = 0xff;
568 controlMessageBuffer[length++] = 0xff;
569 controlMessageBuffer[length++] = 0xff;
570 return controlMessageBuffer;
574 if(clientConnectionPending) {
575 clientConnectionPending =
false;
582 bool DataBlockProtocol::generateResendRequest(
int& length) {
583 length =
static_cast<int>(missingReceiveSegments.size() * (
sizeof(int) +
sizeof(
unsigned short)));
584 if(length +
sizeof(
int) + 1>
sizeof(controlMessageBuffer)) {
589 for(MissingReceiveSegment segment: missingReceiveSegments) {
590 unsigned int segOffset = htonl(static_cast<unsigned int>(segment.offset));
591 unsigned int segLen = htonl(static_cast<unsigned int>(segment.length));
593 memcpy(&controlMessageBuffer[length], &segOffset,
sizeof(segOffset));
594 length +=
sizeof(
unsigned int);
595 memcpy(&controlMessageBuffer[length], &segLen,
sizeof(segLen));
596 length +=
sizeof(
unsigned int);
599 controlMessageBuffer[length++] = RESEND_MESSAGE;
604 void DataBlockProtocol::parseResendMessage(
int length) {
605 missingTransferSegments.clear();
607 int num = length / (
sizeof(
unsigned int) +
sizeof(
unsigned short));
608 int bufferOffset = receiveOffset;
610 for(
int i=0; i<num; i++) {
611 unsigned int segOffsetNet = *
reinterpret_cast<unsigned int*
>(&receiveBuffer[bufferOffset]);
612 bufferOffset +=
sizeof(
unsigned int);
613 unsigned int segLenNet = *
reinterpret_cast<unsigned int*
>(&receiveBuffer[bufferOffset]);
614 bufferOffset +=
sizeof(
unsigned int);
616 int segOffset =
static_cast<int>(ntohl(segOffsetNet));
617 int segLen =
static_cast<int>(ntohl(segLenNet));
619 if(segOffset >= 0 && segLen > 0 && segOffset + segLen < rawValidBytes) {
620 missingTransferSegments.push_back(std::pair<int, int>(
624 LOG_ERROR(
"Requested resend: " << segOffset <<
" - " << (segOffset + segLen));
628 void DataBlockProtocol::parseEofMessage(
int length) {
630 totalReceiveSize =
static_cast<int>(ntohl(*reinterpret_cast<unsigned int*>(
631 &receiveBuffer[receiveOffset])));
632 if(totalReceiveSize < receiveOffset) {
635 if(totalReceiveSize != receiveOffset && receiveOffset != 0) {
637 MissingReceiveSegment missingSeg;
638 missingSeg.offset = receiveOffset;
639 missingSeg.length = totalReceiveSize - receiveOffset;
640 missingSeg.isEof =
true;
641 missingReceiveSegments.push_back(missingSeg);
643 if(missingReceiveSegments.size() > 0) {
644 waitingForMissingSegments =
true;
645 resendMessagePending =
true;
646 receiveOffset = missingReceiveSegments[0].offset;
648 finishedReception =
true;
653 void DataBlockProtocol::resizeReceiveBuffer() {
654 if(totalReceiveSize < 0) {
661 + MAX_OUTSTANDING_BYTES +
sizeof(int);
664 if(static_cast<int>(receiveBuffer.size()) < bufferSize) {
665 receiveBuffer.resize(bufferSize);
Exception class that is used for all protocol exceptions.
void setTransferData(unsigned char *data, int validBytes=0x7FFFFFFF)
Sets the payload data for the next transfer.
const unsigned char * getTransferMessage(int &length)
Gets the next network message for the current transfer.
const unsigned char * getNextControlMessage(int &length)
If a control message is pending to be transmitted, then the message data will be returned by this met...
DataBlockProtocol(bool server, ProtocolType protType, int maxUdpPacketSize)
Creates a new instance.
unsigned char * getReceivedData(int &length)
Returns the data that has been received for the current transfer.
void resetTransfer()
Resets all transfer related internal variables.
unsigned char * getNextReceiveBuffer(int maxLength)
Gets a buffer for receiving the next network message.
int getMaxReceptionSize() const
Returns the maximum payload size that can be received.
bool isConnected() const
Returns true if a remote connection is established.
bool newClientConnected()
Returns true if the last network message has established a new connection from a client.
void processReceivedMessage(int length, bool &transferComplete)
Handles a received network message.
void resetReception(bool dropped)
Resets the message reception.
bool transferComplete()
Returns true if the current transfer has been completed.
void setTransferHeader(unsigned char *data, int headerSize, int transferSize)
Sets a user-defined header that shall be transmitted with the next transfer.
void setTransferValidBytes(int validBytes)
Updates the number of valid bytes in a partial transfer.
unsigned char * getReceivedHeader(int &length)
Returns the header data that has been received for the current transfer.