Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions src/Payload.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,15 +24,23 @@ void Payload::checkFlags(FrameFlags flags) const {
}

std::ostream& operator<<(std::ostream& os, const Payload& payload) {
return os << "[metadata: "
return os
<< "[Metadata("
<< (payload.metadata
? folly::to<std::string>(
payload.metadata->computeChainDataLength())
: "<null>")
<< " data: " << (payload.data
? folly::to<std::string>(
payload.data->computeChainDataLength())
: "<null>")
: "0")
<< (payload.metadata
? "): '" + payload.metadata->cloneAsValue().moveToFbString().substr(0, 80).toStdString() + "'"
: "): <nullptr>")
<< ", Data("
<< (payload.data
? folly::to<std::string>(
payload.data->computeChainDataLength())
: "0")
<< (payload.data
? "): '" + payload.data->cloneAsValue().moveToFbString().substr(0, 80).toStdString() + "'"
: "): <nullptr>")
<< "]";
}

Expand Down
34 changes: 31 additions & 3 deletions src/framing/Frame.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -95,8 +95,34 @@ std::ostream& operator<<(std::ostream& os, ErrorCode errorCode) {
}

std::ostream& operator<<(std::ostream& os, FrameFlags frameFlags) {
std::bitset<16> flags(static_cast<uint16_t>(frameFlags));
return os << flags;
// TODO Match the Flag names with the AllowedFlags in the Frame declarations

std::stringstream ss;
std::string delimeter = "";
if (!!(frameFlags & FrameFlags::NEXT)) {
ss << "NEXT";
delimeter = "|";
}
if (!!(frameFlags & FrameFlags::COMPLETE)) {
ss << delimeter << "COMPLETE";
delimeter = "|";
}
if (!!(frameFlags & FrameFlags::FOLLOWS)) {
ss << delimeter << "FOLLOWS";
delimeter = "|";
}
if (!!(frameFlags & FrameFlags::METADATA)) {
ss << delimeter << "METADATA";
delimeter = "|";
}
if (!!(frameFlags & FrameFlags::IGNORE)) {
ss << delimeter << "IGNORE";
delimeter = "|";
}
if (!delimeter.empty()) {
return os << ss.str();
}
return os << "EMPTY";
}

std::ostream& operator<<(std::ostream& os, const FrameHeader& header) {
Expand DownExpand Up@@ -184,7 +210,9 @@ std::ostream& operator<<(std::ostream& os, const Frame_KEEPALIVE& frame) {
}

std::ostream& operator<<(std::ostream& os, const Frame_SETUP& frame) {
return os << frame.header_ << ", (" << frame.payload_;
return os << frame.header_
<< ", Version: " << frame.versionMajor_ << "." << frame.versionMinor_
<< ", (" << frame.payload_;
}

void Frame_SETUP::moveToSetupPayload(SetupParameters& setupPayload) {
Expand Down
31 changes: 14 additions & 17 deletions src/statemachine/RSocketStateMachine.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -442,7 +442,7 @@ void RSocketStateMachine::handleConnectionFrame(
remoteResumeable_, frame, std::move(payload))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
resumeCache_->resetUpToPosition(frame.position_);
if (mode_ == ReactiveSocketMode::SERVER) {
if (!!(frame.header_.flags_ & FrameFlags::KEEPALIVE_RESPOND)) {
Expand All@@ -464,7 +464,7 @@ void RSocketStateMachine::handleConnectionFrame(
case FrameType::METADATA_PUSH: {
Frame_METADATA_PUSH frame;
if (deserializeFrameOrError(frame, std::move(payload))) {
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
requestResponder_->handleMetadataPush(std::move(frame.metadata_));
}
return;
Expand All@@ -474,7 +474,7 @@ void RSocketStateMachine::handleConnectionFrame(
if (!deserializeFrameOrError(frame, std::move(payload))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
if (resumeCallback_) {
if (resumeCache_->isPositionAvailable(frame.position_)) {
resumeCallback_->onResumeOk();
Expand All@@ -496,7 +496,7 @@ void RSocketStateMachine::handleConnectionFrame(
if (!deserializeFrameOrError(frame, std::move(payload))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;

// TODO: handle INVALID_SETUP, UNSUPPORTED_SETUP, REJECTED_SETUP

Expand DownExpand Up@@ -551,12 +551,12 @@ void RSocketStateMachine::handleStreamFrame(
if (!deserializeFrameOrError(frameRequestN, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frameRequestN;
VLOG(3) << "In:" << frameRequestN;
stateMachine->handleRequestN(frameRequestN.requestN_);
break;
}
case FrameType::CANCEL: {
VLOG(3) << "In:" << Frame_CANCEL();
VLOG(3) << "In:" << Frame_CANCEL();
stateMachine->handleCancel();
break;
}
Expand All@@ -565,7 +565,7 @@ void RSocketStateMachine::handleStreamFrame(
if (!deserializeFrameOrError(framePayload, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << framePayload;
VLOG(3) << "In:" << framePayload;
stateMachine->handlePayload(
std::move(framePayload.payload_),
framePayload.header_.flagsComplete(),
Expand All@@ -577,7 +577,7 @@ void RSocketStateMachine::handleStreamFrame(
if (!deserializeFrameOrError(frameError, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frameError;
VLOG(3) << "In:" << frameError;
stateMachine->handleError(
std::runtime_error(frameError.payload_.moveDataToString()));
break;
Expand DownExpand Up@@ -622,7 +622,7 @@ void RSocketStateMachine::handleUnknownStream(
if (!deserializeFrameOrError(frame, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
auto stateMachine =
streamsFactory_.createChannelResponder(frame.requestN_, streamId);
auto requestSink = requestResponder_->handleRequestChannelCore(
Expand All@@ -635,7 +635,7 @@ void RSocketStateMachine::handleUnknownStream(
if (!deserializeFrameOrError(frame, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
auto stateMachine =
streamsFactory_.createStreamResponder(frame.requestN_, streamId);
requestResponder_->handleRequestStreamCore(
Expand All@@ -647,7 +647,7 @@ void RSocketStateMachine::handleUnknownStream(
if (!deserializeFrameOrError(frame, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
auto stateMachine =
streamsFactory_.createRequestResponseResponder(streamId);
requestResponder_->handleRequestResponseCore(
Expand All@@ -659,7 +659,7 @@ void RSocketStateMachine::handleUnknownStream(
if (!deserializeFrameOrError(frame, std::move(serializedFrame))) {
return;
}
VLOG(3) << "In:" << frame;
VLOG(3) << "In:" << frame;
// no stream tracking is necessary
requestResponder_->handleFireAndForget(
std::move(frame.payload_), streamId);
Expand DownExpand Up@@ -790,15 +790,12 @@ void RSocketStateMachine::requestFireAndForget(Payload request) {
streamsFactory().getNextStreamId(),
FrameFlags::EMPTY,
std::move(std::move(request)));
VLOG(3) << "Out: " << frame;
outputFrameOrEnqueue(frameSerializer_->serializeOut(std::move(frame)));
outputFrameOrEnqueue(std::move(frame));
}

void RSocketStateMachine::metadataPush(std::unique_ptr<folly::IOBuf> metadata) {
Frame_METADATA_PUSH metadataPushFrame(std::move(metadata));
VLOG(3) << "Out: " << metadataPushFrame;
outputFrameOrEnqueue(
frameSerializer_->serializeOut(std::move(metadataPushFrame)));
outputFrameOrEnqueue(std::move(metadataPushFrame));
}

void RSocketStateMachine::outputFrame(std::unique_ptr<folly::IOBuf> frame) {
Expand Down