MCPcopy Create free account
hub / github.com/apple/foundationdb / readMsg

Method readMsg

fdbserver/networktest.actor.cpp:428–464  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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) {

Callers

nothing calls this directly

Calls 8

makeStringFunction · 0.85
delayFunction · 0.85
getMethod · 0.65
readMethod · 0.65
beginMethod · 0.45
endMethod · 0.45
sizeMethod · 0.45
onReadableMethod · 0.45

Tested by

no test coverage detected