This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 453d8eecbb [INLONG-8891][SDK] Optimize compile for dataproxy cpp sdk
(#8893)
453d8eecbb is described below
commit 453d8eecbb47a05ba9b865546c9b1acec3d0f005
Author: doleyzi <[email protected]>
AuthorDate: Tue Sep 12 14:03:39 2023 +0800
[INLONG-8891][SDK] Optimize compile for dataproxy cpp sdk (#8893)
---
.../dataproxy-sdk-cpp/CMakeLists.txt | 16 ++-
.../dataproxy-sdk-cpp/release/inc/inlong_api.h | 4 +-
.../dataproxy-sdk-cpp/release/inc/sdk_conf.h | 2 +-
.../dataproxy-sdk-cpp/release/inc/sdk_msg.h | 6 +-
.../dataproxy-sdk-cpp/src/client/tcp_client.h | 2 +-
.../dataproxy-sdk-cpp/src/config/proxy_info.h | 1 +
.../dataproxy-sdk-cpp/src/config/sdk_conf.cc | 8 +-
.../dataproxy-sdk-cpp/src/core/CMakeLists.txt | 24 +++++
.../dataproxy-sdk-cpp/src/core/api_imp.cc | 10 +-
.../dataproxy-sdk-cpp/src/core/inlong_api.cc | 2 +-
.../dataproxy-sdk-cpp/src/group/CMakeLists.txt | 24 +++++
.../dataproxy-sdk-cpp/src/group/recv_group.cc | 60 +++++-------
.../dataproxy-sdk-cpp/src/group/recv_group.h | 10 +-
.../dataproxy-sdk-cpp/src/group/send_group.cc | 5 +-
.../dataproxy-sdk-cpp/src/group/send_group.h | 2 +-
.../dataproxy-sdk-cpp/src/manager/CMakeLists.txt | 24 +++++
.../dataproxy-sdk-cpp/src/manager/proxy_manager.cc | 24 ++---
.../dataproxy-sdk-cpp/src/manager/proxy_manager.h | 4 +-
.../dataproxy-sdk-cpp/src/manager/send_manager.cc | 8 +-
.../dataproxy-sdk-cpp/src/protocol/CMakeLists.txt | 24 +++++
.../dataproxy-sdk-cpp/src/protocol/msg_protocol.cc | 76 ++++++++++++++
.../dataproxy-sdk-cpp/src/protocol/msg_protocol.h | 88 +++++++++++++++++
.../dataproxy-sdk-cpp/src/utils/capi_constant.h | 9 ++
.../dataproxy-sdk-cpp/src/utils/read_write_mutex.h | 2 +-
.../dataproxy-sdk-cpp/src/utils/send_buffer.h | 109 ++++++++-------------
25 files changed, 392 insertions(+), 152 deletions(-)
diff --git a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/CMakeLists.txt
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/CMakeLists.txt
index 317e4bed6b..bc563f83b8 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/CMakeLists.txt
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/CMakeLists.txt
@@ -30,15 +30,23 @@ include_directories(release/inc)
include_directories(src/base)
include_directories(src/net)
include_directories(src/utils)
+include_directories(src/core)
+include_directories(src/manager)
+include_directories(src/group)
+include_directories(src/protocol)
link_directories(${PROJECT_SOURCE_DIR}/third_party/lib)
link_directories(${PROJECT_SOURCE_DIR}/third_party/lib64)
-add_subdirectory(third_party)
+# add_subdirectory(third_party)
add_subdirectory(src/base)
add_subdirectory(src/net)
add_subdirectory(src/utils)
add_subdirectory(src/config)
+add_subdirectory(src/core)
+add_subdirectory(src/manager)
+add_subdirectory(src/group)
+add_subdirectory(src/protocol)
# add_subdirectory(test)
add_subdirectory(release)
@@ -48,9 +56,9 @@ aux_source_directory(src/base BASE_SRCS)
aux_source_directory(src/net NET_SRCS)
# dynamic library
-# # add_library(dataproxy_sdk SHARED ${BASE_SRCS} ${NET_SRCS})
-# # set_target_properties(dataproxy_sdk PROPERTIES PREFIX "")
-# # target_link_libraries(dataproxy_sdk -llog4cplus -lsnappy -lcurl)
+# add_library(dataproxy_sdk SHARED ${BASE_SRCS} ${NET_SRCS})
+# set_target_properties(dataproxy_sdk PROPERTIES PREFIX "")
+# target_link_libraries(dataproxy_sdk -llog4cplus -lsnappy -lcurl)
# static library
add_library(dataproxy_sdk STATIC ${BASE_SRCS} ${NET_SRCS})
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/inlong_api.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/inlong_api.h
index 23ca465359..458155d9bc 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/inlong_api.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/inlong_api.h
@@ -32,7 +32,7 @@ using UserCallBack =
std::function<int32_t(const char *, const char *, const char *, int32_t,
const int64_t, const char *)>;
-class InLongApiImp;
+class ApiImp;
class InLongApi {
public:
@@ -49,7 +49,7 @@ public:
int32_t CloseApi(int32_t max_waitms);
private:
- std::shared_ptr<InLongApiImp> api_impl_;
+ std::shared_ptr<ApiImp> api_impl_;
};
} // namespace inlong
#endif // INLONG_SDK_API_H
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_conf.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_conf.h
index d31407dc34..b230f7a7bc 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_conf.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_conf.h
@@ -73,7 +73,7 @@ public:
std::string manager_cluster_url_;
uint32_t manager_update_interval_; // Automatic update interval, minutes
uint32_t manager_url_timeout_; // URL parsing timeout, seconds
- uint32_t max_tcp_num_;
+ uint32_t max_proxy_num_;
uint32_t msg_type_;
// Network parameters
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_msg.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_msg.h
index a8cae707fa..1114eb17db 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_msg.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/release/inc/sdk_msg.h
@@ -25,7 +25,7 @@
namespace inlong {
using UserCallBack = std::function<int32_t(const char*, const char*, const
char*, int32_t, const int64_t, const char*)>;
-struct UserMsg {
+struct SdkMsg {
std::string msg_;
std::string client_ip_;
int64_t report_time_;
@@ -35,7 +35,7 @@ struct UserMsg {
std::string user_client_ip_;
std::string data_pack_format_attr_;
- UserMsg(const std::string& mmsg, const std::string& mclient_ip, int64_t
mreport_time, UserCallBack mcb,
+ SdkMsg(const std::string& mmsg, const std::string& mclient_ip, int64_t
mreport_time, UserCallBack mcb,
const std::string& attr, const std::string& u_ip, int64_t u_time)
: msg_(mmsg),
client_ip_(mclient_ip),
@@ -45,7 +45,7 @@ struct UserMsg {
user_client_ip_(u_ip),
data_pack_format_attr_(attr){}
};
-using UserMsgPtr = std::shared_ptr<UserMsg>;
+using SdkMsgPtr = std::shared_ptr<SdkMsg>;
} // namespace inlong
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/client/tcp_client.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/client/tcp_client.h
index 5e665813dd..be0b6d915a 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/client/tcp_client.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/client/tcp_client.h
@@ -24,7 +24,7 @@
#include "../utils/capi_constant.h"
#include "../utils/read_write_mutex.h"
#include "../utils/send_buffer.h"
-#include "../utils/stat.h"
+#include "stat.h"
#include <queue>
namespace inlong {
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/proxy_info.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/proxy_info.h
index b5fccb5213..d433e3a295 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/proxy_info.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/proxy_info.h
@@ -18,6 +18,7 @@
*/
#include <string>
+#include <vector>
#ifndef INLONG_SDK_PROXY_INFO_H
#define INLONG_SDK_PROXY_INFO_H
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/sdk_conf.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/sdk_conf.cc
index fcb7d7da19..b5b613717b 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/sdk_conf.cc
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/config/sdk_conf.cc
@@ -102,7 +102,7 @@ void SdkConfig::defaultInit() {
manager_cluster_url_ = constants::kBusClusterURL;
manager_update_interval_ = constants::kBusUpdateInterval;
manager_url_timeout_ = constants::kBusURLTimeout;
- max_tcp_num_ = constants::kMaxBusNum;
+ max_proxy_num_ = constants::kMaxBusNum;
local_ip_ = constants::kSerIP;
local_port_ = constants::kSerPort;
@@ -296,9 +296,9 @@ void SdkConfig::InitBusParam(const rapidjson::Value &doc) {
if (doc.HasMember("max_tcp_num") && doc["max_tcp_num"].IsInt() &&
doc["max_tcp_num"].GetInt() > 0) {
const rapidjson::Value &obj = doc["max_tcp_num"];
- max_tcp_num_ = obj.GetInt();
+ max_proxy_num_ = obj.GetInt();
} else {
- max_tcp_num_ = constants::kMaxBusNum;
+ max_proxy_num_ = constants::kMaxBusNum;
}
if (doc.HasMember("group_ids") && doc["group_ids"].IsString()) {
@@ -403,7 +403,7 @@ void SdkConfig::ShowClientConfig() {
: "false");
LOG_INFO("manager_update_interval: minutes" << manager_update_interval_);
LOG_INFO("manager_url_timeout: " << manager_url_timeout_);
- LOG_INFO("max_tcp_num: " << max_tcp_num_);
+ LOG_INFO("max_tcp_num: " << max_proxy_num_);
LOG_INFO("msg_type: " << msg_type_);
LOG_INFO("enable_tcp_nagle: " << enable_tcp_nagle_ ? "true" : "false");
LOG_INFO("enable_setaffinity: " << enable_setaffinity_ ? "true" : "false");
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/CMakeLists.txt
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/CMakeLists.txt
new file mode 100644
index 0000000000..51b34e5983
--- /dev/null
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/CMakeLists.txt
@@ -0,0 +1,24 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+cmake_minimum_required(VERSION 3.1)
+
+aux_source_directory(. CORE_SRCS)
+
+add_library(inlong_core STATIC ${CORE_SRCS})
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/api_imp.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/api_imp.cc
index b80f1cc40f..b6c8760f7a 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/api_imp.cc
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/api_imp.cc
@@ -37,7 +37,7 @@ int32_t ApiImp::InitApi(const char *config_file_path) {
}
max_msg_length_ = std::min(SdkConfig::getInstance()->max_msg_size_,
SdkConfig::getInstance()->pack_size_);
- local_ip_ = SdkConfig::getInstance()->ser_ip_;
+ local_ip_ = SdkConfig::getInstance()->local_ip_;
return DoInit();
}
@@ -95,7 +95,7 @@ int32_t ApiImp::CloseApi(int32_t max_waitms) {
int32_t ApiImp::DoInit() {
LOG_INFO(
- "tdbus sdk cpp start Init, version:" << constants::kTDBusCAPIVersion);
+ "tdbus sdk cpp start Init, version:" << constants::kVersion);
signal(SIGPIPE, SIG_IGN);
@@ -103,11 +103,11 @@ int32_t ApiImp::DoInit() {
ProxyManager::GetInstance()->Init();
- for (int i = 0; i < SdkConfig::getInstance()->inlong_group_ids_.size(); i++)
{
+ for (int i = 0; i < SdkConfig::getInstance()->group_ids_.size(); i++) {
LOG_INFO("DoInit CheckBidConf inlong_group_id:"
- << SdkConfig::getInstance()->inlong_group_ids_[i]);
+ << SdkConfig::getInstance()->group_ids_[i]);
ProxyManager::GetInstance()->CheckBidConf(
- SdkConfig::getInstance()->inlong_group_ids_[i], false);
+ SdkConfig::getInstance()->group_ids_[i], false);
}
return InitManager();
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/inlong_api.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/inlong_api.cc
index 0a8b75fa07..cef90fdb46 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/inlong_api.cc
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/core/inlong_api.cc
@@ -19,7 +19,7 @@
#include "../core/api_imp.h"
namespace inlong {
-InLongApi::InLongApi() { api_impl_ = std::make_shared<InLongApiImp>(); };
+InLongApi::InLongApi() { api_impl_ = std::make_shared<ApiImp>(); };
InLongApi::~InLongApi() { api_impl_->CloseApi(10); }
int32_t InLongApi::InitApi(const char *config_path) {
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/CMakeLists.txt
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/CMakeLists.txt
new file mode 100644
index 0000000000..2967f2c651
--- /dev/null
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/CMakeLists.txt
@@ -0,0 +1,24 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+cmake_minimum_required(VERSION 3.1)
+
+aux_source_directory(. GROUP_SRCS)
+
+add_library(inlong_group STATIC ${GROUP_SRCS})
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.cc
index 5d9075c7c0..4e58fe9be4 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.cc
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.cc
@@ -19,7 +19,7 @@
#include "../utils/utils.h"
#include "api_code.h"
-#include "msg_protocol.h"
+#include "../protocol/msg_protocol.h"
#include <cstdlib>
#include <functional>
@@ -28,7 +28,7 @@ const uint32_t ATTR_LENGTH = 10;
const uint32_t DEFAULT_PACK_ATTR = 400;
RecvGroup::RecvGroup(const std::string &inlong_group_id, const std::string
&inlong_stream_id,
std::shared_ptr<SendManager> send_manager)
- : cur_len_(0), inlong_group_id_(inlong_group_id),
inlong_stream_id_(inlong_stream_id), groupId_num_(0), tid_num_(0),
+ : cur_len_(0), inlong_group_id_(inlong_group_id),
inlong_stream_id_(inlong_stream_id), groupId_num_(0), streamId_num_(0),
msg_type_(SdkConfig::getInstance()->msg_type_),
data_capacity_(SdkConfig::getInstance()->buf_size_),
send_manager_(send_manager) {
@@ -69,7 +69,7 @@ int32_t RecvGroup::SendData(const std::string &msg, const
std::string &groupId,
int32_t RecvGroup::DoDispatchMsg() {
last_pack_time_ = Utils::getCurrentMsTime();
std::lock_guard<std::mutex> lck(mutex_);
- if (groupId_.empty()) {
+ if (inlong_group_id_.empty()) {
LOG_ERROR("groupId is empty, check!!");
return SdkCode::kInvalidInput;
}
@@ -77,7 +77,7 @@ int32_t RecvGroup::DoDispatchMsg() {
LOG_ERROR("no msg in msg_set, check!");
return SdkCode::kFailGetRevGroup;
}
- auto send_group = send_manager_->GetSendGroup(groupId_);
+ auto send_group = send_manager_->GetSendGroup(inlong_group_id_);
if (send_group == nullptr) {
LOG_ERROR("failed to get send_buf, something gets wrong, checkout!");
return SdkCode::kFailGetSendBuf;
@@ -90,16 +90,16 @@ int32_t RecvGroup::DoDispatchMsg() {
}
uint32_t total_length = 0;
- std::vector<UserMsgPtr> msgs_to_dispatch;
+ std::vector<SdkMsgPtr> msgs_to_dispatch;
while (!msgs_.empty()) {
- UserMsgPtr msg = msgs_.front();
- if (msg->msg.size() + total_length + ATTR_LENGTH >
+ SdkMsgPtr msg = msgs_.front();
+ if (msg->msg_.size() + total_length + ATTR_LENGTH >
SdkConfig::getInstance()->pack_size_) {
break;
}
msgs_to_dispatch.push_back(msg);
msgs_.pop();
- total_length = msg->msg.size() + total_length + ATTR_LENGTH;
+ total_length = msg->msg_.size() + total_length + ATTR_LENGTH;
}
cur_len_ = cur_len_ - total_length;
@@ -137,7 +137,7 @@ void RecvGroup::AddMsg(const std::string &msg, std::string
client_ip,
std::string data_pack_format_attr =
"__addcol1__reptime=" + Utils::getFormatTime(data_time_) +
"&__addcol2__ip=" + client_ip;
- msgs_.push(std::make_shared<UserMsg>(msg, client_ip, data_time_, call_back,
+ msgs_.push(std::make_shared<SdkMsg>(msg, client_ip, data_time_, call_back,
data_pack_format_attr, user_client_ip,
user_report_time));
@@ -152,7 +152,7 @@ bool RecvGroup::ShouldPack(int32_t msg_len) {
return false;
}
-bool RecvGroup::PackMsg(std::vector<UserMsgPtr> &msgs, char *pack_data,
+bool RecvGroup::PackMsg(std::vector<SdkMsgPtr> &msgs, char *pack_data,
uint32_t &out_len, uint32_t uniq_id) {
if (pack_data == nullptr) {
LOG_ERROR("nullptr, failed to allocate memory for buf");
@@ -161,20 +161,20 @@ bool RecvGroup::PackMsg(std::vector<UserMsgPtr> &msgs,
char *pack_data,
uint32_t idx = 0;
for (auto &it : msgs) {
if (msg_type_ >= 5) {
- *(uint32_t *)(&pack_buf_[idx]) = htonl(it->msg.size());
+ *(uint32_t *)(&pack_buf_[idx]) = htonl(it->msg_.size());
idx += sizeof(uint32_t);
}
- memcpy(&pack_buf_[idx], it->msg.data(), it->msg.size());
- idx += static_cast<uint32_t>(it->msg.size());
+ memcpy(&pack_buf_[idx], it->msg_.data(), it->msg_.size());
+ idx += static_cast<uint32_t>(it->msg_.size());
// add attrlen|attr
if (SdkConfig::getInstance()->isAttrDataPackFormat()) {
- *(uint32_t *)(&pack_buf_[idx]) = htonl(it->data_pack_format_attr.size());
+ *(uint32_t *)(&pack_buf_[idx]) =
htonl(it->data_pack_format_attr_.size());
idx += sizeof(uint32_t);
- memcpy(&pack_buf_[idx], it->data_pack_format_attr.data(),
- it->data_pack_format_attr.size());
- idx += static_cast<uint32_t>(it->data_pack_format_attr.size());
+ memcpy(&pack_buf_[idx], it->data_pack_format_attr_.data(),
+ it->data_pack_format_attr_.size());
+ idx += static_cast<uint32_t>(it->data_pack_format_attr_.size());
}
if (msg_type_ == 2 || msg_type_ == 3) {
@@ -215,7 +215,7 @@ bool RecvGroup::PackMsg(std::vector<UserMsgPtr> &msgs, char
*pack_data,
uint32_t char_groupId_flag = 0;
std::string groupId_streamId_char;
uint16_t groupId_num = 0, streamId_num = 0;
- if (SdkConfig::getInstance()->enableChargroupId() || groupId_num_ == 0 ||
+ if (SdkConfig::getInstance()->enableChar() || groupId_num_ == 0 ||
streamId_num_ == 0) {
groupId_num = 0;
streamId_num = 0;
@@ -232,10 +232,10 @@ bool RecvGroup::PackMsg(std::vector<UserMsgPtr> &msgs,
char *pack_data,
std::string attr;
if (SdkConfig::getInstance()->enableTraceIP()) {
if (groupId_streamId_char.empty())
- attr = "node1ip=" + SdkConfig::getInstance()->ser_ip_ +
+ attr = "node1ip=" + SdkConfig::getInstance()->local_ip_ +
"&rtime1=" + std::to_string(Utils::getCurrentMsTime());
else
- attr = groupId_streamId_char + "&node1ip=" +
SdkConfig::getInstance()->ser_ip_ +
+ attr = groupId_streamId_char + "&node1ip=" +
SdkConfig::getInstance()->local_ip_ +
"&rtime1=" + std::to_string(Utils::getCurrentMsTime());
} else {
attr = topic_desc_;
@@ -298,12 +298,6 @@ bool RecvGroup::PackMsg(std::vector<UserMsgPtr> &msgs,
char *pack_data,
attr += "&cp=snappy";
attr += "&cnt=" + std::to_string(cnt);
attr += "&sid=" + std::string(Utils::getSnowflakeId());
- if (SdkConfig::getInstance()
- ->is_from_DC_) {
//&__addcol1_reptime=yyyymmddHHMMSS&__addcol2__ip=BBB&f=dc
- attr += "&__addcol1_reptime=" +
- Utils::getFormatTime(Utils::getCurrentMsTime()) +
- "&__addcol2__ip=" + SdkConfig::getInstance()->ser_ip_ + "&f=dc";
- }
*(uint32_t *)bodyBegin = htonl(attr.size());
bodyBegin += sizeof(uint32_t);
@@ -339,7 +333,7 @@ void RecvGroup::DispatchMsg(bool exit) {
}
}
std::shared_ptr<SendBuffer>
-RecvGroup::BuildSendBuf(std::vector<UserMsgPtr> &msgs) {
+RecvGroup::BuildSendBuf(std::vector<SdkMsgPtr> &msgs) {
if (msgs.empty()) {
LOG_ERROR("pack msgs is empty.");
return nullptr;
@@ -361,8 +355,8 @@ RecvGroup::BuildSendBuf(std::vector<UserMsgPtr> &msgs) {
}
send_buffer->setLen(len);
send_buffer->setMsgCnt(msg_cnt);
- send_buffer->setgroupId(groupId_);
- send_buffer->setstreamId(streamId_);
+ send_buffer->setInlongGroupId(inlong_group_id_);
+ send_buffer->setStreamId(inlong_stream_id_);
send_buffer->setUniqId(uniq_id);
send_buffer->setIsPacked(true);
for (auto it : msgs) {
@@ -372,11 +366,11 @@ RecvGroup::BuildSendBuf(std::vector<UserMsgPtr> &msgs) {
return send_buffer;
}
-void RecvGroup::CallbalkToUsr(std::vector<UserMsgPtr> &msgs) {
+void RecvGroup::CallbalkToUsr(std::vector<SdkMsgPtr> &msgs) {
for (auto &it : msgs) {
- if (it->cb) {
- it->cb(groupId_.data(), streamId_.data(), it->msg.data(), it->msg.size(),
- it->user_report_time, it->user_client_ip.data());
+ if (it->cb_) {
+ it->cb_(inlong_group_id_.data(), inlong_stream_id_.data(),
it->msg_.data(), it->msg_.size(),
+ it->user_report_time_, it->user_client_ip_.data());
}
}
}
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.h
index ea85660573..3624e26c60 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/recv_group.h
@@ -28,13 +28,13 @@
#include "../utils/atomic.h"
#include "../utils/noncopyable.h"
#include "sdk_conf.h"
-#include "user_msg.h"
+
namespace inlong {
class RecvGroup {
private:
char *pack_buf_;
- std::queue<UserMsgPtr> msgs_;
+ std::queue<SdkMsgPtr> msgs_;
uint32_t data_capacity_;
uint32_t cur_len_;
AtomicInt pack_err_;
@@ -67,7 +67,7 @@ public:
int32_t SendData(const std::string &msg, const std::string &groupId,
const std::string &streamId, const std::string &client_ip,
uint64_t report_time, UserCallBack call_back);
- bool PackMsg(std::vector<UserMsgPtr> &msgs, char *pack_data,
+ bool PackMsg(std::vector<SdkMsgPtr> &msgs, char *pack_data,
uint32_t &out_len, uint32_t uniq_id);
void DispatchMsg(
bool
@@ -77,8 +77,8 @@ public:
std::string groupId() const { return inlong_group_id_; }
- std::shared_ptr<SendBuffer> BuildSendBuf(std::vector<UserMsgPtr> &msgs);
- void CallbalkToUsr(std::vector<UserMsgPtr> &msgs);
+ std::shared_ptr<SendBuffer> BuildSendBuf(std::vector<SdkMsgPtr> &msgs);
+ void CallbalkToUsr(std::vector<SdkMsgPtr> &msgs);
};
using RecvGroupPtr = std::shared_ptr<RecvGroup>;
} // namespace inlong
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.cc
index 319b2a9ad5..57f1c574b4 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.cc
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.cc
@@ -19,7 +19,7 @@
#include "send_group.h"
#include "api_code.h"
-#include "proxy_conf_manager.h"
+#include "proxy_manager.h"
#include <algorithm>
#include <random>
@@ -121,8 +121,7 @@ void SendGroup::UpdateConf(std::error_code error) {
ClearOldTcpClients();
ProxyInfoVec new_proxy_info;
- if (proxyConfManager::GetInstance()->GetproxyByBid(
- group_id_, new_proxy_info) != kSuccess ||
+ if (ProxyManager::GetInstance()->GetProxy(group_id_, new_proxy_info) !=
kSuccess ||
new_proxy_info.empty()) {
update_conf_timer_->expires_after(std::chrono::milliseconds(kTimerMinute));
update_conf_timer_->async_wait(
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.h
index 9b929ca4d8..b1659e3631 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/group/send_group.h
@@ -22,7 +22,7 @@
#include "../config/proxy_info.h"
#include "../utils/send_buffer.h"
-#include "tcp_client.h"
+#include "../client/tcp_client.h"
#include <queue>
namespace inlong {
const int kTimerMiSeconds = 10;
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/CMakeLists.txt
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/CMakeLists.txt
new file mode 100644
index 0000000000..e3348fbfa8
--- /dev/null
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/CMakeLists.txt
@@ -0,0 +1,24 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+cmake_minimum_required(VERSION 3.1)
+
+aux_source_directory(. MANAGER_SRCS)
+
+add_library(inlong_manager STATIC ${MANAGER_SRCS})
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.cc
index 5548a441ba..28edc7e413 100644
---
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.cc
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.cc
@@ -39,7 +39,7 @@ ProxyManager::~ProxyManager() {
cond_.notify_one();
}
void ProxyManager::Init() {
- timeout_ = SdkConfig::getInstance()->bus_URL_timeout_;
+ timeout_ = SdkConfig::getInstance()->manager_url_timeout_;
if (__sync_bool_compare_and_swap(&inited_, false, true)) {
update_conf_thread_ = std::thread(&ProxyManager::Update, this);
}
@@ -50,7 +50,7 @@ void ProxyManager::Update() {
std::unique_lock<std::mutex> con_lck(cond_mutex_);
if (cond_.wait_for(con_lck,
std::chrono::minutes(
- SdkConfig::getInstance()->bus_update_interval_),
+ SdkConfig::getInstance()->manager_update_interval_),
[this]() { return update_flag_; })) {
if (exit_flag_)
break;
@@ -74,19 +74,19 @@ void ProxyManager::DoUpdate() {
}
{
- unique_read_lock<read_write_mutex> rdlck(bid_2_cluster_rwmutex_);
+ unique_read_lock<read_write_mutex> rdlck(groupid_2_cluster_rwmutex_);
for (auto &groupid2cluster : groupid_2_cluster_map_) {
std::string url;
- if (SdkConfig::getInstance().enable_proxy_URL_from_cluster_)
- url = SdkConfig::getInstance().proxy_cluster_URL_;
+ if (SdkConfig::getInstance()->enable_manager_url_from_cluster_)
+ url = SdkConfig::getInstance()->manager_cluster_url_;
else {
- url = SdkConfig::getInstance().proxy_URL_ + "/" +
groupid2cluster.first;
+ url = SdkConfig::getInstance()->manager_url_ + "/" +
groupid2cluster.first;
}
- std::string post_data = "ip=" + SdkConfig::getInstance().ser_ip_ +
- "&version=" + constants::kTDBusCAPIVersion +
+ std::string post_data = "ip=" + SdkConfig::getInstance()->local_ip_ +
+ "&version=" + constants::kVersion +
"&protocolType=" + constants::kProtocolType;
- LOG_WARN("get inlong_group_id:%s proxy cfg url:%s, post_data:%s",
- groupid2cluster.first.c_str(), url.c_str(), post_data.c_str());
+ LOG_WARN("get inlong_group_id:%s proxy cfg url:%s, post_data:%s"<<
+ groupid2cluster.first.c_str()<<url.c_str()<< post_data.c_str());
std::string meta_data;
int32_t ret;
@@ -185,8 +185,8 @@ int32_t ProxyManager::ParseAndGet(const std::string
&groupid,
return SdkCode::kSuccess;
}
-int32_t ProxyManager::GetBusByBid(const std::string &groupid,
- BusInfoVec &proxy_info_vec) {
+int32_t ProxyManager::GetProxy(const std::string &groupid,
+ ProxyInfoVec &proxy_info_vec) {
unique_read_lock<read_write_mutex> rdlck(groupid_2_proxy_map_rwmutex_);
auto it = groupid_2_proxy_map_.find(groupid);
if (it == groupid_2_proxy_map_.end()) {
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.h
index 1c76f98457..50a27e478b 100644
---
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.h
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/proxy_manager.h
@@ -34,7 +34,7 @@ private:
read_write_mutex groupid_2_cluster_rwmutex_;
read_write_mutex groupid_2_proxy_map_rwmutex_;
- std::unordered_map<std::string, int32_t> groupid_2_cluster_map_cluster_map_;
+ std::unordered_map<std::string, int32_t> groupid_2_cluster_map_;
std::unordered_map<std::string, ProxyInfoVec> groupid_2_proxy_map_;
bool update_flag_;
std::mutex cond_mutex_;
@@ -56,7 +56,7 @@ public:
void Update();
void DoUpdate();
void Init();
- int32_t GetBusByBid(const std::string &groupid, ProxyInfoVec
&proxy_info_vec);
+ int32_t GetProxy(const std::string &groupid, ProxyInfoVec &proxy_info_vec);
bool IsBusExist(const std::string &groupid);
};
} // namespace inlong
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/send_manager.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/send_manager.cc
index a5cbda53d1..3c45537319 100644
---
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/send_manager.cc
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/manager/send_manager.cc
@@ -25,7 +25,7 @@ SendManager::SendManager() : send_group_idx_(0) {
for (int32_t i = 0; i < SdkConfig::getInstance()->group_ids_.size(); i++) {
LOG_INFO("SendManager, group_id:"
<< SdkConfig::getInstance()->group_ids_[i] << " send group num:"
- << SdkConfig::getInstance()->per_group_id_thread_nums_);
+ << SdkConfig::getInstance()->per_groupid_thread_nums_);
DoAddSendGroup(SdkConfig::getInstance()->group_ids_[i]);
}
}
@@ -39,7 +39,7 @@ SendGroupPtr SendManager::GetSendGroup(const std::string
&group_id) {
}
bool SendManager::AddSendGroup(const std::string &group_id) {
- if (!BusConfManager::GetInstance()->IsBusExist(group_id)) {
+ if (!ProxyManager::GetInstance()->IsBusExist(group_id)) {
LOG_ERROR("bus is not exist." << group_id);
return false;
}
@@ -55,8 +55,8 @@ void SendManager::DoAddSendGroup(const std::string &group_id)
{
return;
}
std::vector<SendGroupPtr> send_group;
- send_group.reserve(SdkConfig::getInstance()->per_group_id_thread_nums_);
- for (int32_t j = 0; j < SdkConfig::getInstance()->per_group_id_thread_nums_;
+ send_group.reserve(SdkConfig::getInstance()->per_groupid_thread_nums_);
+ for (int32_t j = 0; j < SdkConfig::getInstance()->per_groupid_thread_nums_;
j++) {
send_group.push_back(std::make_shared<SendGroup>(group_id));
}
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/CMakeLists.txt
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/CMakeLists.txt
new file mode 100644
index 0000000000..00b15aac3a
--- /dev/null
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/CMakeLists.txt
@@ -0,0 +1,24 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+cmake_minimum_required(VERSION 3.1)
+
+aux_source_directory(. PROTOCOL_SRCS)
+
+add_library(inlong_protocol STATIC ${PROTOCOL_SRCS})
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.cc
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.cc
new file mode 100644
index 0000000000..8d561a78ec
--- /dev/null
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.cc
@@ -0,0 +1,76 @@
+/**
+* Licensed to the Apache Software Foundation (ASF) under one
+* or more contributor license agreements. See the NOTICE file
+* distributed with this work for additional information
+* regarding copyright ownership. The ASF licenses this file
+* to you under the Apache License, Version 2.0 (the
+* "License"); you may not use this file except in compliance
+* with the License. You may obtain a copy of the License at
+*
+* http://www.apache.org/licenses/LICENSE-2.0
+*
+* Unless required by applicable law or agreed to in writing,
+* software distributed under the License is distributed on an
+* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+* KIND, either express or implied. See the License for the
+* specific language governing permissions and limitations
+* under the License.
+*/
+
+#include "msg_protocol.h"
+
+#include "../utils/logger.h"
+#include <arpa/inet.h>
+#include <assert.h>
+#include <snappy.h>
+#include <string.h>
+namespace inlong {
+char APIEncode::recvBuf[APIEncode::kRecvLen] = {0};
+
+void APIEncode::decodeProtocoMsg(SendBuffer *buf) {
+ memset(recvBuf, 0x0, kRecvLen);
+ // //LOG_DEBUG("print buf content, %s", buf->content());
+ memcpy(recvBuf, buf->content(), buf->len());
+ // memcpy(recvBuf, buf->content(), buf->len());
+
+ char *p = recvBuf;
+
+ uint32_t total_len = 0;
+ memcpy(&total_len, p, 4);
+ total_len = ntohl(total_len);
+ p += 4;
+
+ uint8_t msg_type = 0;
+ memcpy(&msg_type, p, 1);
+ ++p;
+
+ uint32_t body_len = 0;
+ memcpy(&body_len, p, 4);
+ body_len = ntohl(body_len);
+ p += 4;
+
+ std::string body_uncompress;
+ bool res = snappy::Uncompress(p, body_len, &body_uncompress);
+ // LOG_DEBUG("uncompress res:%d, body_len:%d, body:%s", res,
+ // body_uncompress.size(), body_uncompress.c_str());
+
+ p += body_len; // skip body
+
+ uint32_t attr_len = 0;
+ memcpy(&attr_len, p, 4);
+ attr_len = ntohl(attr_len);
+ p += 4;
+ // LOG_WARN("decode,attr_len:%d", attr_len);
+
+ char attr[attr_len + 1];
+ memset(attr, 0x0, attr_len);
+ memcpy(attr, p, attr_len);
+
+ if (attr_len <= 0) {
+ LOG_ERROR("decode protoco error!!!!!!!, total_len:"
+ << total_len << " msg_type:" << msg_type
+ << " body_len:" << body_len << "attr_len :" << attr_len);
+ return;
+ }
+}
+} // namespace inlong
\ No newline at end of file
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.h
new file mode 100644
index 0000000000..ab1cfcb806
--- /dev/null
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/protocol/msg_protocol.h
@@ -0,0 +1,88 @@
+/**
+* Licensed to the Apache Software Foundation (ASF) under one
+* or more contributor license agreements. See the NOTICE file
+* distributed with this work for additional information
+* regarding copyright ownership. The ASF licenses this file
+* to you under the Apache License, Version 2.0 (the
+* "License"); you may not use this file except in compliance
+* with the License. You may obtain a copy of the License at
+*
+* http://www.apache.org/licenses/LICENSE-2.0
+*
+* Unless required by applicable law or agreed to in writing,
+* software distributed under the License is distributed on an
+* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+* KIND, either express or implied. See the License for the
+* specific language governing permissions and limitations
+* under the License.
+*/
+
+#ifndef INLONG_SDK_MSG_PROTOCOL_H
+#define INLONG_SDK_MSG_PROTOCOL_H
+
+#include "../utils/send_buffer.h"
+#include <memory>
+#include <stdint.h>
+#include <string>
+#include <vector>
+
+namespace inlong {
+// totalLen(4)|msgtype(1)|bodyLen(4)|body(x)|attrLen(4)|attr(attr_len)
+#pragma pack(1)
+struct ProtocolMsgHead {
+ uint32_t total_len;
+ char msg_type;
+};
+struct ProtocolMsgBody {
+ uint32_t body_len;
+ char body[0];
+};
+struct ProtocolMsgTail {
+ uint32_t attr_len;
+ char attr[0];
+};
+//
totalLen(4)|msgtype(1)|bid_num(2)|tid_num(2)|ext_field(2)|data_time(4)|cnt(2)|uniq(4)|bodyLen(4)|body(x)|attrLen(4)|attr(attr_len)|magic(2)
+struct BinaryMsgHead {
+ uint32_t total_len;
+ char msg_type;
+ uint16_t bid_num;
+ uint16_t tid_num;
+ uint16_t ext_field;
+ uint32_t data_time;
+ uint16_t cnt;
+ uint32_t uniq;
+};
+
+struct BinaryMsgAck {
+ uint32_t total_len;
+ char msg_type;
+ uint32_t uniq;
+ uint16_t attr_len;
+ char attr[0];
+ uint16_t magic;
+};
+
+struct BinaryHB {
+ uint32_t total_len;
+ char msg_type;
+ uint32_t data_time;
+ uint8_t body_ver; // body_ver=1
+ uint32_t body_len; // body_len=0
+ char body[0];
+ uint16_t attr_len; // attr_len=0;
+ char attr[0];
+ uint16_t magic;
+};
+#pragma pack()
+
+class APIEncode {
+public:
+ static void decodeProtocoMsg(SendBuffer *buf);
+
+ static const uint32_t kRecvLen = 1024 * 1024 * 11;
+ static char recvBuf[kRecvLen];
+};
+
+} // namespace inlong
+
+#endif // INLONG_SDK_MSG_PROTOCOL_H
\ No newline at end of file
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/capi_constant.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/capi_constant.h
index 2f191909ba..e045171de3 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/capi_constant.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/capi_constant.h
@@ -75,6 +75,15 @@ static const bool kEnableSetAffinity = false;
static const uint32_t kMaskCPUAffinity = 0xff;
static const uint16_t kExtendField = 0;
+// http basic auth
+static const char kBasicAuthHeader[] = "Authorization:";
+static const char kBasicAuthPrefix[] = "Basic";
+static const char kBasicAuthSeparator[] = " ";
+static const char kBasicAuthJoiner[] = ":";
+static const char kProtocolType [] = "TCP";
+
+static const uint32_t kMaxAttrLen = 2048;
+
} // namespace constants
} // namespace inlong
#endif // INLONG_SDK_CONSTANT_H
\ No newline at end of file
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/read_write_mutex.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/read_write_mutex.h
index 44bb7471fe..ef73450d49 100644
---
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/read_write_mutex.h
+++
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/read_write_mutex.h
@@ -23,7 +23,7 @@
#include <condition_variable>
#include <mutex>
-namespace dataproxy_sdk {
+namespace inlong {
// wirte operation add lock:unique_read_lock<read_write_mutex> lock( rwmutex );
// read operation add lock:unique_write_lock<read_write_mutex> lock(rwmutex);
diff --git
a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/send_buffer.h
b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/send_buffer.h
index 4aa2f55cf2..16bde14941 100644
--- a/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/send_buffer.h
+++ b/inlong-sdk/dataproxy-sdk-twins/dataproxy-sdk-cpp/src/utils/send_buffer.h
@@ -23,41 +23,39 @@
#include <mutex>
#include <string>
-#include "sdk_core.h"
-// #include "executor_thread_pool.h"
+#include "atomic.h"
#include "logger.h"
#include "noncopyable.h"
-// #include "socket_connection.h"
-
-namespace dataproxy_sdk {
+#include "sdk_msg.h"
+#include <asio.hpp>
+#include <deque>
+#include <queue>
+
+namespace inlong {
+class Connection;
+using SteadyTimerPtr = std::shared_ptr<asio::steady_timer>;
+using ConnectionPtr = std::shared_ptr<Connection>;
class SendBuffer : noncopyable {
private:
uint32_t uniq_id_;
- bool is_used_;
- bool is_packed_; // is packed completed
- char *content_; // send data
- uint32_t size_; // buf_size
+ std::atomic<bool> is_used_;
+ std::atomic<bool> is_packed_;
+ char *content_;
+ uint32_t size_;
int32_t msg_cnt_;
- uint32_t len_; // send data len
+ uint32_t len_;
std::string inlong_group_id_;
std::string inlong_stream_id_;
AtomicInt already_send_;
- uint64_t first_send_time_; // ms
- uint64_t latest_send_time_; // ms
- ConnectionPtr target_; // send conn
- std::vector<UserMsgPtr> user_msg_set_;
-
-public:
- std::mutex mutex_;
- SteadyTimerPtr timeout_timer_; // timeout, resend
- AtomicInt fail_create_conn_; // create conn fail count
+ uint64_t first_send_time_;
+ uint64_t latest_send_time_;
+ std::vector<SdkMsgPtr> user_msg_vector_;
public:
SendBuffer(uint32_t size)
- : is_used_(false), is_packed_(false), mutex_(), msg_cnt_(0), len_(0),
- inlong_group_id_(), inlong_stream_id_(), first_send_time_(0),
- latest_send_time_(0), target_(nullptr), uniq_id_(0),
- timeout_timer_(nullptr), size_(size) {
+ : uniq_id_(0), is_used_(false), is_packed_(false), size_(size),
+ msg_cnt_(0), len_(0), inlong_group_id_(), inlong_stream_id_(),
first_send_time_(0),
+ latest_send_time_(0) {
content_ = new char[size];
if (content_) {
memset(content_, 0x0, size);
@@ -67,7 +65,6 @@ public:
if (content_) {
delete[] content_;
}
- content_ = nullptr;
}
char *content() { return content_; }
@@ -75,47 +72,27 @@ public:
void setMsgCnt(const int32_t &msg_cnt) { msg_cnt_ = msg_cnt; }
uint32_t len() { return len_; }
void setLen(const uint32_t len) { len_ = len; }
- std::string inlong_group_id() { return inlong_group_id_; }
- std::string inlong_stream_id() { return inlong_stream_id_; }
- void setGroupid(const std::string &inlong_group_id) {
- inlong_group_id_ = inlong_group_id;
- }
- void setStreamid(const std::string &inlong_stream_id) {
- inlong_stream_id_ = inlong_stream_id;
- }
- uint64_t firstSendTime() const { return first_send_time_; }
- void setFirstSendTime(const uint64_t &first_send_time) {
- first_send_time_ = first_send_time;
- }
- uint64_t latestSendTime() const { return latest_send_time_; }
- void setLatestSendTime(const uint64_t &latest_send_time) {
- latest_send_time_ = latest_send_time;
- }
+ std::string bid() { return inlong_group_id_; }
+ std::string tid() { return inlong_stream_id_; }
+ void setInlongGroupId(const std::string &inlong_group_id) { inlong_group_id_
= inlong_group_id; }
+ void setStreamId(const std::string &inlong_stream_id) { inlong_stream_id_ =
inlong_stream_id; }
- ConnectionPtr target() const { return target_; }
- void setTarget(ConnectionPtr &target) { target_ = target; }
-
- inline void increaseRetryNum() { already_send_.increment(); }
- inline int32_t getAlreadySend() { return already_send_.get(); }
-
- uint32_t uniqId() const { return uniq_id_; }
void setUniqId(const uint32_t &uniq_id) { uniq_id_ = uniq_id; }
- void addUserMsg(UserMsgPtr u_msg) { user_msg_set_.push_back(u_msg); }
+ void addUserMsg(SdkMsgPtr msg) { user_msg_vector_.push_back(msg); }
void doUserCallBack() {
- LOG_TRACE("failed to send msg, start user call_back");
- for (auto it : user_msg_set_) {
- if (it->cb) {
- it->cb(inlong_group_id_.data(), inlong_stream_id_.data(),
- it->msg.data(), it->msg.size(), it->user_report_time,
- it->user_client_ip.data());
+ for (auto it : user_msg_vector_) {
+ if (it->cb_) {
+ it->cb_(inlong_group_id_.data(), inlong_stream_id_.data(),
it->msg_.data(), it->msg_.size(),
+ it->user_report_time_, it->user_client_ip_.data());
}
}
}
- void reset() {
- uint32_t record_uid = uniq_id_; // for debug
-
+ void releaseBuf() {
+ if (!is_used_) {
+ return;
+ }
uniq_id_ = 0;
is_used_ = false;
is_packed_ = false;
@@ -127,23 +104,15 @@ public:
already_send_.getAndSet(0);
first_send_time_ = 0;
latest_send_time_ = 0;
- target_ = nullptr;
- if (timeout_timer_) {
- timeout_timer_->cancel();
- timeout_timer_ = nullptr;
- }
- user_msg_set_.clear();
+ user_msg_vector_.clear();
+ AtomicInt fail_create_conn_;
fail_create_conn_.getAndSet(0);
- LOG_TRACE("reset senfbuf(uid:%d) successfully", record_uid);
}
- bool isPacked() const { return is_packed_; }
void setIsPacked(bool is_packed) { is_packed_ = is_packed; }
-
- bool isUsed() const { return is_used_; }
- void setIsUsed(bool is_used) { is_used_ = is_used; }
};
-
-} // namespace dataproxy_sdk
+typedef std::shared_ptr<SendBuffer> SendBufferPtrT;
+typedef std::queue<SendBufferPtrT> SendBufferPtrDeque;
+} // namespace inlong
#endif // INLONG_SDK_SEND_BUFFER_H
\ No newline at end of file