This is an automated email from the ASF dual-hosted git repository.
phrocker pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/nifi-minifi-cpp.git
The following commit(s) were added to refs/heads/master by this push:
new 7514caa MINIFICPP-706 - RawSiteToSite: remove code duplication
7514caa is described below
commit 7514caacbf4722f7e3c47224ec62d918d6f4804d
Author: Arpad Boda <[email protected]>
AuthorDate: Wed Jan 9 14:15:13 2019 +0100
MINIFICPP-706 - RawSiteToSite: remove code duplication
This closes #470.
Signed-off-by: Marc Parisi <[email protected]>
---
libminifi/include/sitetosite/RawSocketProtocol.h | 11 +--
libminifi/include/sitetosite/SiteToSiteClient.h | 8 +-
libminifi/src/sitetosite/RawSocketProtocol.cpp | 98 +++++-------------------
libminifi/src/sitetosite/SiteToSiteClient.cpp | 37 ---------
4 files changed, 23 insertions(+), 131 deletions(-)
diff --git a/libminifi/include/sitetosite/RawSocketProtocol.h
b/libminifi/include/sitetosite/RawSocketProtocol.h
index 395b597..1624ed6 100644
--- a/libminifi/include/sitetosite/RawSocketProtocol.h
+++ b/libminifi/include/sitetosite/RawSocketProtocol.h
@@ -112,9 +112,6 @@ class RawSiteToSiteClient : public
sitetosite::SiteToSiteClient {
}
- void setPeer(std::unique_ptr<SiteToSitePeer> peer) {
- peer_ = std::move(peer);
- }
/**
* Provides a reference to the time out
* @returns timeout
@@ -146,18 +143,14 @@ class RawSiteToSiteClient : public
sitetosite::SiteToSiteClient {
virtual int writeRequestType(RequestType type);
// read Request Type
virtual int readRequestType(RequestType &type);
+
// read Respond
virtual int readRespond(const std::shared_ptr<Transaction> &transaction,
RespondCode &code, std::string &message);
// write respond
virtual int writeRespond(const std::shared_ptr<Transaction> &transaction,
RespondCode code, std::string message);
// getRespondCodeContext
virtual RespondCodeContext *getRespondCodeContext(RespondCode code) {
- for (unsigned int i = 0; i < sizeof(SiteToSiteRequest::respondCodeContext)
/ sizeof(RespondCodeContext); i++) {
- if (SiteToSiteRequest::respondCodeContext[i].code == code) {
- return &SiteToSiteRequest::respondCodeContext[i];
- }
- }
- return NULL;
+ return SiteToSiteClient::getRespondCodeContext(code);
}
// Creation of a new transaction, return the transaction ID if success,
diff --git a/libminifi/include/sitetosite/SiteToSiteClient.h
b/libminifi/include/sitetosite/SiteToSiteClient.h
index 6469e82..e1fffc3 100644
--- a/libminifi/include/sitetosite/SiteToSiteClient.h
+++ b/libminifi/include/sitetosite/SiteToSiteClient.h
@@ -221,12 +221,8 @@ class SiteToSiteClient : public core::Connectable {
// deleteTransaction
virtual void deleteTransaction(std::string transactionID);
- virtual void tearDown();
-
- // write Request Type
- virtual int writeRequestType(RequestType type);
- // read Request Type
- virtual int readRequestType(RequestType &type);
+ virtual void tearDown() = 0;
+
// read Respond
virtual int readResponse(const std::shared_ptr<Transaction> &transaction,
RespondCode &code, std::string &message);
// write respond
diff --git a/libminifi/src/sitetosite/RawSocketProtocol.cpp
b/libminifi/src/sitetosite/RawSocketProtocol.cpp
index c0bf499..316397a 100644
--- a/libminifi/src/sitetosite/RawSocketProtocol.cpp
+++ b/libminifi/src/sitetosite/RawSocketProtocol.cpp
@@ -395,97 +395,37 @@ bool
RawSiteToSiteClient::getPeerList(std::vector<PeerStatus> &peers) {
}
}
-int RawSiteToSiteClient::writeRequestType(RequestType type) {
- if (type >= MAX_REQUEST_TYPE)
- return -1;
-
- return peer_->writeUTF(SiteToSiteRequest::RequestTypeStr[type]);
-}
-
-int RawSiteToSiteClient::readRequestType(RequestType &type) {
- std::string requestTypeStr;
-
- int ret = peer_->readUTF(requestTypeStr);
-
- if (ret <= 0)
- return ret;
+ int RawSiteToSiteClient::writeRequestType(RequestType type) {
+ if (type >= MAX_REQUEST_TYPE)
+ return -1;
- for (int i = NEGOTIATE_FLOWFILE_CODEC; i <= SHUTDOWN; i++) {
- if (SiteToSiteRequest::RequestTypeStr[i] == requestTypeStr) {
- type = (RequestType) i;
- return ret;
- }
+ return peer_->writeUTF(SiteToSiteRequest::RequestTypeStr[type]);
}
- return -1;
-}
-
-int RawSiteToSiteClient::readRespond(const std::shared_ptr<Transaction>
&transaction, RespondCode &code, std::string &message) {
- uint8_t firstByte;
-
- int ret = peer_->read(firstByte);
-
- if (ret <= 0 || firstByte != CODE_SEQUENCE_VALUE_1)
- return -1;
-
- uint8_t secondByte;
-
- ret = peer_->read(secondByte);
-
- if (ret <= 0 || secondByte != CODE_SEQUENCE_VALUE_2)
- return -1;
-
- uint8_t thirdByte;
-
- ret = peer_->read(thirdByte);
-
- if (ret <= 0)
- return ret;
+ int RawSiteToSiteClient::readRequestType(RequestType &type) {
+ std::string requestTypeStr;
- code = (RespondCode) thirdByte;
+ int ret = peer_->readUTF(requestTypeStr);
- RespondCodeContext *resCode = getRespondCodeContext(code);
-
- if (resCode == NULL) {
- // Not a valid respond code
- return -1;
- }
- if (resCode->hasDescription) {
- ret = peer_->readUTF(message);
if (ret <= 0)
- return -1;
- }
- return 3 + message.size();
-}
+ return ret;
-int RawSiteToSiteClient::writeRespond(const std::shared_ptr<Transaction>
&transaction, RespondCode code, std::string message) {
- RespondCodeContext *resCode = getRespondCodeContext(code);
+ for (int i = NEGOTIATE_FLOWFILE_CODEC; i <= SHUTDOWN; i++) {
+ if (SiteToSiteRequest::RequestTypeStr[i] == requestTypeStr) {
+ type = (RequestType) i;
+ return ret;
+ }
+ }
- if (resCode == NULL) {
- // Not a valid respond code
return -1;
}
- uint8_t codeSeq[3];
- codeSeq[0] = CODE_SEQUENCE_VALUE_1;
- codeSeq[1] = CODE_SEQUENCE_VALUE_2;
- codeSeq[2] = (uint8_t) code;
-
- int ret = peer_->write(codeSeq, 3);
-
- if (ret != 3)
- return -1;
+int RawSiteToSiteClient::readRespond(const std::shared_ptr<Transaction>
&transaction, RespondCode &code, std::string &message) {
+ return readResponse(transaction, code, message);
+}
- if (resCode->hasDescription) {
- ret = peer_->writeUTF(message);
- if (ret > 0) {
- return (3 + ret);
- } else {
- return ret;
- }
- } else {
- return 3;
- }
+int RawSiteToSiteClient::writeRespond(const std::shared_ptr<Transaction>
&transaction, RespondCode code, std::string message) {
+ return writeResponse(transaction, code, message);
}
bool RawSiteToSiteClient::negotiateCodec() {
diff --git a/libminifi/src/sitetosite/SiteToSiteClient.cpp
b/libminifi/src/sitetosite/SiteToSiteClient.cpp
index 61fefdf..3bc91fe 100644
--- a/libminifi/src/sitetosite/SiteToSiteClient.cpp
+++ b/libminifi/src/sitetosite/SiteToSiteClient.cpp
@@ -25,31 +25,6 @@ namespace nifi {
namespace minifi {
namespace sitetosite {
-int SiteToSiteClient::writeRequestType(RequestType type) {
- if (type >= MAX_REQUEST_TYPE)
- return -1;
-
- return peer_->writeUTF(SiteToSiteRequest::RequestTypeStr[type]);
-}
-
-int SiteToSiteClient::readRequestType(RequestType &type) {
- std::string requestTypeStr;
-
- int ret = peer_->readUTF(requestTypeStr);
-
- if (ret <= 0)
- return ret;
-
- for (int i = NEGOTIATE_FLOWFILE_CODEC; i <= SHUTDOWN; i++) {
- if (SiteToSiteRequest::RequestTypeStr[i] == requestTypeStr) {
- type = (RequestType) i;
- return ret;
- }
- }
-
- return -1;
-}
-
int SiteToSiteClient::readResponse(const std::shared_ptr<Transaction>
&transaction, RespondCode &code, std::string &message) {
uint8_t firstByte;
@@ -133,18 +108,6 @@ int SiteToSiteClient::writeResponse(const
std::shared_ptr<Transaction> &transact
}
}
-void SiteToSiteClient::tearDown() {
- if (peer_state_ >= ESTABLISHED) {
- logger_->log_debug("Site2Site Protocol tearDown");
- // need to write shutdown request
- writeRequestType(SHUTDOWN);
- }
-
- known_transactions_.clear();
- peer_->Close();
- peer_state_ = IDLE;
-}
-
bool SiteToSiteClient::transferFlowFiles(const
std::shared_ptr<core::ProcessContext> &context, const
std::shared_ptr<core::ProcessSession> &session) {
std::shared_ptr<FlowFileRecord> flow =
std::static_pointer_cast<FlowFileRecord>(session->get());