| 28 | } |
| 29 | |
| 30 | void ReceiverWorker::doRead() |
| 31 | { |
| 32 | // are there enough bytes to read something |
| 33 | bool can_read_more = true; |
| 34 | |
| 35 | while (m_socket.bytesAvailable() > 0 || can_read_more) |
| 36 | { |
| 37 | |
| 38 | if (m_socket.bytesAvailable() > 0) |
| 39 | can_read_more = true; |
| 40 | |
| 41 | m_buffer.append(m_socket.readAll()); |
| 42 | |
| 43 | // read the size of the next field if haven't already |
| 44 | if (!m_state.size_read) |
| 45 | { |
| 46 | |
| 47 | if (m_buffer.size() < m_state.bytes_read + FIELD_SIZE_NBYTES) |
| 48 | { |
| 49 | /// can't read, need to wait for more bytes |
| 50 | can_read_more = false; |
| 51 | continue; |
| 52 | } |
| 53 | |
| 54 | /// enough bytes to read size |
| 55 | m_state.msg_size = ArrayToInt(m_buffer.mid(m_state.bytes_read, 4)); |
| 56 | m_state.bytes_read += 4; |
| 57 | m_state.size_read = true; |
| 58 | } |
| 59 | else |
| 60 | { |
| 61 | |
| 62 | if (m_buffer.size() < m_state.bytes_read + m_state.msg_size) |
| 63 | { |
| 64 | /// can't read, need to wait for more bytes |
| 65 | can_read_more = false; |
| 66 | continue; |
| 67 | } |
| 68 | |
| 69 | marshalling.deserialize(m_buffer.data() + m_state.bytes_read, m_state.msg_size); |
| 70 | |
| 71 | auto msg = marshalling.get_msg(); |
| 72 | handleMessage(msg); |
| 73 | |
| 74 | m_state.bytes_read += m_state.msg_size; |
| 75 | m_state.size_read = false; |
| 76 | |
| 77 | m_state.msg_processed++; |
| 78 | |
| 79 | /// reset the buffer every MSG_PER_BUFFER messages |
| 80 | if (m_state.msg_processed == MSG_PER_BUFFER) |
| 81 | { |
| 82 | m_state.msg_processed = 0; |
| 83 | m_buffer.remove(0, m_state.bytes_read); |
| 84 | m_state.bytes_read = 0; |
| 85 | } |
| 86 | } |
| 87 | } |
nothing calls this directly
no test coverage detected