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


Reply via email to