| 426 | NetworkAddress randomRemote() { return remotes[nondeterministicRandom()->randomInt(0, remotes.size())]; } |
| 427 | |
| 428 | ACTOR static Future<Standalone<StringRef>> readMsg(P2PNetworkTest* self, Reference<IConnection> conn) { |
| 429 | state Standalone<StringRef> buffer = makeString(sizeof(int)); |
| 430 | state int writeOffset = 0; |
| 431 | state bool gotHeader = false; |
| 432 | |
| 433 | // Fill buffer sequentially until the initial bytesToRead is read (or more), then read |
| 434 | // intended message size and add it to bytesToRead, continue if needed until bytesToRead is 0. |
| 435 | loop { |
| 436 | int stutter = self->waitReadMilliseconds.get(); |
| 437 | if (stutter > 0) { |
| 438 | wait(delay(stutter / 1e3)); |
| 439 | } |
| 440 | |
| 441 | int len = conn->read((uint8_t*)buffer.begin() + writeOffset, (uint8_t*)buffer.end()); |
| 442 | writeOffset += len; |
| 443 | self->bytesReceived += len; |
| 444 | |
| 445 | // If buffer is complete, either process it as a header or return it |
| 446 | if (writeOffset == buffer.size()) { |
| 447 | if (gotHeader) { |
| 448 | return buffer; |
| 449 | } else { |
| 450 | gotHeader = true; |
| 451 | int msgSize = *(int*)buffer.begin(); |
| 452 | if (msgSize == 0) { |
| 453 | return Standalone<StringRef>(); |
| 454 | } |
| 455 | buffer = makeString(msgSize); |
| 456 | writeOffset = 0; |
| 457 | } |
| 458 | } |
| 459 | |
| 460 | if (len == 0) { |
| 461 | wait(conn->onReadable()); |
| 462 | wait(delay(0, TaskPriority::ReadSocket)); |
| 463 | } |
| 464 | } |
| 465 | } |
| 466 | |
| 467 | ACTOR static Future<Void> writeMsg(P2PNetworkTest* self, Reference<IConnection> conn, StringRef msg) { |
nothing calls this directly
no test coverage detected