This is an automated email from the ASF dual-hosted git repository.
gosonzhang 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 8a11a89e4 [INLONG-7466][TubeMQ] Adjust code style issues (#7467)
8a11a89e4 is described below
commit 8a11a89e44d5c7522bb360ddcba9214113599dda
Author: Goson Zhang <[email protected]>
AuthorDate: Tue Feb 28 16:52:29 2023 +0800
[INLONG-7466][TubeMQ] Adjust code style issues (#7467)
---
.../tubemq-client-cpp/example/producer/test_producer.cc | 2 +-
.../tubemq-client-cpp/src/baseproducer.cc | 5 +++--
.../tubemq-client-cpp/src/const_config.h | 2 +-
.../tubemq-client-go/client/heartbeat.go | 6 ++++--
.../tubemq-client-twins/tubemq-client-go/errs/errs.go | 2 +-
.../tubemq-client-go/metadata/metadata.go | 2 +-
.../src/python/example/test_producer.py | 2 ++
.../tubemq-client-python/src/python/tubemq/client.py | 3 ++-
.../inlong/tubemq/server/broker/BrokerServiceServer.java | 6 +++---
.../tubemq/server/broker/web/BrokerAdminServlet.java | 7 ++-----
.../org/apache/inlong/tubemq/server/master/TMaster.java | 16 ++++++++--------
11 files changed, 28 insertions(+), 25 deletions(-)
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/example/producer/test_producer.cc
b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/example/producer/test_producer.cc
index cef7ee35c..8ece52634 100644
---
a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/example/producer/test_producer.cc
+++
b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/example/producer/test_producer.cc
@@ -140,7 +140,7 @@ int main(int argc, char* argv[]) {
}
}
- while (MessageSentCallback::kTotalCounter.Get() < (long)msg_count) {
+ while (MessageSentCallback::kTotalCounter.Get() < (int64)msg_count) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
auto stop = std::chrono::steady_clock::now();
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/baseproducer.cc
b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/baseproducer.cc
index 14df0e766..284e738b3 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/baseproducer.cc
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/baseproducer.cc
@@ -93,7 +93,8 @@ bool BaseProducer::Start(string& err_info, const
ProducerConfig& config) {
status_.CompareAndSet(tb_config::kMasterRegistering,
tb_config::kMasterUnRegistered);
return false;
}
- status_.CompareAndSet(tb_config::kMasterRegistering,
tb_config::kMasterRegistered); // register2Master done, change status_ to `2`
+ // register2Master done, change status_ to `2`
+ status_.CompareAndSet(tb_config::kMasterRegistering,
tb_config::kMasterRegistered);
// set heartbeat timer
heart_beat_timer_ = TubeMQService::Instance()->CreateTimer();
@@ -678,7 +679,7 @@ bool BaseProducer::processHBResponseM2P(int32_t&
error_code, string& err_info,
error_code = rsp_protocol->code_;
err_info = rsp_protocol->error_msg_;
return false;
- };
+ }
HeartResponseM2P rsp_m2p;
bool result = rsp_m2p.ParseFromArray(rsp_protocol->rsp_body_.data().c_str(),
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/const_config.h
b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/const_config.h
index cbe87b5ca..9a891a34c 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/const_config.h
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-cpp/src/const_config.h
@@ -118,7 +118,7 @@ enum RegisterMasterStatus {
kMasterUnRegistered = 0,
kMasterRegistering = 1,
kMasterRegistered = 2
-};
+};
// enum MasterHBStatus
enum MasterHBStatus {
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/heartbeat.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/heartbeat.go
index 7bd9bfa82..0aee413bd 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/heartbeat.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/heartbeat.go
@@ -133,11 +133,13 @@ func (h *heartbeatManager) consumerHB2Master() {
go func() {
err := h.consumer.register2Master(!hbNoNode)
if err != nil {
- log.Warnf("[CONSUMER] heartBeat2Master failure
to (%s) : %s, client=%s", h.consumer.master.Address, rsp.GetErrMsg(),
h.consumer.clientID)
+ log.Warnf("[CONSUMER] heartBeat2Master failure
to (%s) : %s, client=%s",
+ h.consumer.master.Address,
rsp.GetErrMsg(), h.consumer.clientID)
return
}
h.registerMaster(h.consumer.master.Address)
- log.Infof("[CONSUMER] heartBeat2Master success to (%s),
client=%s", h.consumer.master.Address, h.consumer.clientID)
+ log.Infof("[CONSUMER] heartBeat2Master success to (%s),
client=%s",
+ h.consumer.master.Address, h.consumer.clientID)
}()
}
}
diff --git a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/errs/errs.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/errs/errs.go
index ac8167e74..0d80fa4bc 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/errs/errs.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/errs/errs.go
@@ -80,7 +80,7 @@ type Error struct {
Msg string
}
-// Error() implements the Error interface.
+// Error implements the Error interface.
func (e *Error) Error() string {
return fmt.Sprintf("code: %d, msg:%s", e.Code, e.Msg)
}
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/metadata/metadata.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/metadata/metadata.go
index 5ee20c0e2..03183f898 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/metadata/metadata.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/metadata/metadata.go
@@ -56,7 +56,7 @@ func (m *Metadata) SetSubscribeInfo(sub *SubscribeInfo) {
m.subscribeInfo = sub
}
-// ReadStatus sets the status.
+// SetReadStatus sets the status.
func (m *Metadata) SetReadStatus(status int32) {
m.readStatus = status
}
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/example/test_producer.py
b/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/example/test_producer.py
index e15d78cbf..87fb6eacd 100644
---
a/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/example/test_producer.py
+++
b/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/example/test_producer.py
@@ -31,6 +31,7 @@ kSuccessCounter = 0
kFailCounter = 0
counter_lock = Lock()
+
# Reference: java producer: MixedUtils.buildTestData, only for demo
def build_test_data(msg_data_size):
transmit_data = "This is a test data!"
@@ -41,6 +42,7 @@ def build_test_data(msg_data_size):
data += transmit_data[:msg_data_size - len(data)]
return data
+
def send_callback(error_code):
global counter_lock
global kTotalCounter
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/tubemq/client.py
b/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/tubemq/client.py
index 3a977b2b9..64ec52bcb 100644
---
a/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/tubemq/client.py
+++
b/inlong-tubemq/tubemq-client-twins/tubemq-client-python/src/python/tubemq/client.py
@@ -86,7 +86,8 @@ class Producer(tubemq_client.TubeMQProducer):
if not result:
print("StopTubeMQService failure, reason is:" + err_info)
exit(1)
-
+
+
class Consumer(tubemq_client.TubeMQConsumer):
def __init__(self,
master_addr,
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/BrokerServiceServer.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/BrokerServiceServer.java
index 94a901804..388801a51 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/BrokerServiceServer.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/BrokerServiceServer.java
@@ -612,7 +612,7 @@ public class BrokerServiceServer implements
BrokerReadService, BrokerWriteServic
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
// get and check clientId field
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuffer,
result)) {
builder.setErrCode(result.getErrCode());
@@ -808,7 +808,7 @@ public class BrokerServiceServer implements
BrokerReadService, BrokerWriteServic
public RegisterResponseB2C consumerRegisterC2B(RegisterRequestC2B request,
final String rmtAddress,
boolean overtls) throws Throwable {
- ProcessResult result = new ProcessResult();
+ final ProcessResult result = new ProcessResult();
RegisterResponseB2C.Builder builder = RegisterResponseB2C.newBuilder();
builder.setSuccess(false);
builder.setCurrOffset(-1);
@@ -822,7 +822,7 @@ public class BrokerServiceServer implements
BrokerReadService, BrokerWriteServic
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
final StringBuilder strBuffer = new StringBuilder(512);
// get and check clientId field
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuffer,
result)) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/web/BrokerAdminServlet.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/web/BrokerAdminServlet.java
index 2a73b0e67..d88df4d14 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/web/BrokerAdminServlet.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/broker/web/BrokerAdminServlet.java
@@ -1070,9 +1070,6 @@ public class BrokerAdminServlet extends
AbstractWebHandler {
return;
}
// get the history offset in the time range
- // read history data
- int totalCnt = 0;
- // locate start offset
int maxRetryCnt = 50;
long requestOffset = msgStore.getStartOffsetByTimeStamp(recStartTime);
if (!getStoredGroupHisOffsets(groupName, msgStore,
@@ -1099,7 +1096,7 @@ public class BrokerAdminServlet extends
AbstractWebHandler {
// after
Map<String, Map<String, Map<Integer, GroupOffsetInfo>>>
aftGroupOffsetMap =
getGroupOffsetInfo(WebFieldDef.COMPSGROUPNAME, groupNameSet,
topicSet);
- Map<String, Map<Integer, GroupOffsetInfo>> aftTopicPartMap =
+ final Map<String, Map<Integer, GroupOffsetInfo>> aftTopicPartMap =
aftGroupOffsetMap.get(groupName);
// build result
WebParameterUtils.buildSuccessWithDataRetBegin(sBuffer);
@@ -1142,7 +1139,7 @@ public class BrokerAdminServlet extends
AbstractWebHandler {
sBuffer.append("],\"partCount\":").append(partCnt).append("}");
}
sBuffer.append("]}");
- WebParameterUtils.buildSuccessWithDataRetEnd(sBuffer, totalCnt);
+ WebParameterUtils.buildSuccessWithDataRetEnd(sBuffer, 1);
}
/**
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/TMaster.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/TMaster.java
index 9cb7ced13..37707ebfc 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/TMaster.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/TMaster.java
@@ -319,7 +319,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());
@@ -400,7 +400,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());
@@ -537,7 +537,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());
@@ -730,7 +730,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());
@@ -933,7 +933,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
final String rmtAddress,
boolean overtls) throws Exception {
// #lizard forgives
- ProcessResult result = new ProcessResult();
+ final ProcessResult result = new ProcessResult();
final StringBuilder strBuff = new StringBuilder(512);
RegisterResponseM2B.Builder builder = RegisterResponseM2B.newBuilder();
builder.setSuccess(false);
@@ -1043,7 +1043,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
final String rmtAddress,
boolean overtls) throws Exception {
// #lizard forgives
- ProcessResult result = new ProcessResult();
+ final ProcessResult result = new ProcessResult();
final StringBuilder strBuff = new StringBuilder(512);
// set response field
HeartResponseM2B.Builder builder = HeartResponseM2B.newBuilder();
@@ -1199,7 +1199,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());
@@ -1370,7 +1370,7 @@ public class TMaster extends HasThread implements
MasterService, Stoppable {
builder.setErrMsg(result.getErrMsg());
return builder.build();
}
- CertifiedInfo certifiedInfo = (CertifiedInfo) result.getRetData();
+ final CertifiedInfo certifiedInfo = (CertifiedInfo)
result.getRetData();
if (!PBParameterUtils.checkClientId(request.getClientId(), strBuff,
result)) {
builder.setErrCode(result.getErrCode());
builder.setErrMsg(result.getErrMsg());