From 31036381a172b17a965aa74ab161ea833ffe49bd Mon Sep 17 00:00:00 2001 From: xxxx Date: Mon, 25 Jul 2022 09:33:11 +0800 Subject: [PATCH] Description:latency count Match-id-8c8f6d8c16f660408afc7cf0bfe5b5c2f6dbbfbb --- common/include/constants_dinput.h | 6 + common/include/dinput_errcode.h | 2 + .../common/include/dinput_softbus_define.h | 2 + .../distributed_input_sink_transport.h | 2 + .../src/distributed_input_sink_transport.cpp | 40 ++++++ .../distributed_input_source_transport.h | 15 +++ .../distributed_input_source_transport.cpp | 125 ++++++++++++++++-- utils/include/dinput_utils_tool.h | 1 + utils/src/dinput_utils_tool.cpp | 10 ++ 9 files changed, 190 insertions(+), 13 deletions(-) diff --git a/common/include/constants_dinput.h b/common/include/constants_dinput.h index 67a7a33..945c5c4 100644 --- a/common/include/constants_dinput.h +++ b/common/include/constants_dinput.h @@ -86,6 +86,10 @@ namespace DistributedInput { constexpr int32_t INPUT_LOAD_SA_TIMEOUT_MS = 10000; + constexpr int32_t INPUT_LATENCY_DELAYTIME_US = 50 * 1000; + + constexpr int32_t INPUT_LATENCY_DELAY_TIMES = 60; + constexpr int32_t SESSION_WAIT_TIMEOUT_SECOND = 5; /* The input device is a keyboard or has buttons. */ @@ -135,6 +139,8 @@ namespace DistributedInput { const std::string DH_ID_PREFIX = "Input_"; + const std::string DINPUT_SPLIT_COMMA = ", "; + enum class EHandlerMsgType { DINPUT_SINK_EVENT_HANDLER_MSG = 1, DINPUT_SOURCE_EVENT_HANDLER_MSG = 2 diff --git a/common/include/dinput_errcode.h b/common/include/dinput_errcode.h index 64905e6..348a9a6 100644 --- a/common/include/dinput_errcode.h +++ b/common/include/dinput_errcode.h @@ -51,6 +51,7 @@ namespace DistributedInput { constexpr int32_t ERR_DH_INPUT_SERVER_SINK_TRANSPORT_RESPSTART_FAIL = -64007; constexpr int32_t ERR_DH_INPUT_SERVER_SINK_TRANSPORT_RESPSTOP_FAIL = -64008; constexpr int32_t ERR_DH_INPUT_SERVER_SINK_TRANSPORT_SENDMESSAGE_FAIL = -64009; + constexpr int32_t ERR_DH_INPUT_SERVER_SINK_TRANSPORT_RESP_LATENCY_FAIL = -64010; // service source error code constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_INJECT_REGISTER_FAIL = -65000; constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_INJECT_UNREGISTER_FAIL = -65001; @@ -86,6 +87,7 @@ namespace DistributedInput { constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_MANAGER_UNPREPARE_MSG_IS_BAD = -65031; constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_MANAGER_START_MSG_IS_BAD = -65032; constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_MANAGER_STOP_MSG_IS_BAD = -65033; + constexpr int32_t ERR_DH_INPUT_SERVER_SOURCE_TRANSPORT_LATENCY_FAIL = -65034; // handler error code constexpr int32_t ERR_DH_INPUT_SINK_HANDLER_INIT_SINK_SA_FAIL = -66000; constexpr int32_t ERR_DH_INPUT_SINK_HANDLER_INIT_SOURCE_SA_FAIL = -66001; diff --git a/services/common/include/dinput_softbus_define.h b/services/common/include/dinput_softbus_define.h index 7931cd0..8bc65d8 100644 --- a/services/common/include/dinput_softbus_define.h +++ b/services/common/include/dinput_softbus_define.h @@ -60,10 +60,12 @@ namespace DistributedInput { const uint32_t TRANS_SINK_MSG_ONSTART = 3; const uint32_t TRANS_SINK_MSG_ONSTOP = 4; const uint32_t TRANS_SINK_MSG_BODY_DATA = 5; + const uint32_t TRANS_SINK_MSG_LATENCY = 6; const uint32_t TRANS_SOURCE_MSG_PREPARE = 11; const uint32_t TRANS_SOURCE_MSG_UNPREPARE = 12; const uint32_t TRANS_SOURCE_MSG_START = 13; const uint32_t TRANS_SOURCE_MSG_STOP = 14; + const uint32_t TRANS_SOURCE_MSG_LATENCY = 15; } // namespace DistributedInput } // namespace DistributedHardware } // namespace OHOS diff --git a/services/sink/transport/include/distributed_input_sink_transport.h b/services/sink/transport/include/distributed_input_sink_transport.h index e6eb250..964a064 100644 --- a/services/sink/transport/include/distributed_input_sink_transport.h +++ b/services/sink/transport/include/distributed_input_sink_transport.h @@ -45,6 +45,7 @@ public: int32_t RespUnprepareRemoteInput(const int32_t sessionId, std::string &smsg); int32_t RespStartRemoteInput(const int32_t sessionId, std::string &smsg); int32_t RespStopRemoteInput(const int32_t sessionId, std::string &smsg); + int32_t RespLatency(const int32_t sessionId, std::string &smsg); int32_t OnSessionOpened(int32_t sessionId, int32_t result); void OnSessionClosed(int32_t sessionId); @@ -69,6 +70,7 @@ private: void NotifyUnprepareRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyStartRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyStopRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); + void NotifyLatency(int32_t sessionId, const nlohmann::json &recMsg); private: std::string deviceId_; diff --git a/services/sink/transport/src/distributed_input_sink_transport.cpp b/services/sink/transport/src/distributed_input_sink_transport.cpp index 71439ec..25d4cce 100644 --- a/services/sink/transport/src/distributed_input_sink_transport.cpp +++ b/services/sink/transport/src/distributed_input_sink_transport.cpp @@ -225,6 +225,23 @@ int32_t DistributedInputSinkTransport::RespStopRemoteInput( } } +int32_t DistributedInputSinkTransport::RespLatency(const int32_t sessionId, std::string &smsg) +{ + if (sessionId <= 0) { + DHLOGE("RespLatency error, sessionId <= 0."); + return ERR_DH_INPUT_SERVER_SINK_TRANSPORT_RESP_LATENCY_FAIL; + } + + int32_t ret = SendMessage(sessionId, smsg); + if (ret != DH_SUCCESS) { + DHLOGE("RespLatency error, SendMessage fail."); + return ERR_DH_INPUT_SERVER_SINK_TRANSPORT_RESP_LATENCY_FAIL; + } + + DHLOGI("RespLatency sessionId:%s, smsg:%s.", GetAnonyInt32(sessionId).c_str(), smsg.c_str()); + return DH_SUCCESS; +} + int32_t DistributedInputSinkTransport::SendMessage(int32_t sessionId, std::string &message) { DHLOGI("start SendMessage"); @@ -381,6 +398,25 @@ void DistributedInputSinkTransport::NotifyStopRemoteInput(int32_t sessionId, con callback_->onStopRemoteInput(sessionId, inputTypes); } +void DistributedInputSinkTransport::NotifyLatency(int32_t sessionId, const nlohmann::json &recMsg) +{ + if (recMsg.contains(DINPUT_SOFTBUS_KEY_DEVICE_ID) != true) { + DHLOGE("OnBytesReceived message is error, not contain deviceId."); + return; + } + + if (!recMsg[DINPUT_SOFTBUS_KEY_DEVICE_ID].is_string()) { + DHLOGE("OnBytesReceived cmdType TRANS_SOURCE_MSG_LATENCY, data type is error."); + return; + } + + nlohmann::json jsonStr; + jsonStr[DINPUT_SOFTBUS_KEY_CMD_TYPE] = TRANS_SINK_MSG_LATENCY; + jsonStr[DINPUT_SOFTBUS_KEY_RESP_VALUE] = true; + std::string smsg = jsonStr.dump(); + RespLatency(sessionId, smsg); +} + void DistributedInputSinkTransport::HandleSessionData(int32_t sessionId, const std::string& message) { if (callback_ == nullptr) { @@ -422,6 +458,10 @@ void DistributedInputSinkTransport::HandleSessionData(int32_t sessionId, const s NotifyStopRemoteInput(sessionId, recMsg); break; } + case TRANS_SOURCE_MSG_LATENCY: { + NotifyLatency(sessionId, recMsg); + break; + } default: DHLOGE("OnBytesReceived cmdType is undefined."); break; diff --git a/services/source/transport/include/distributed_input_source_transport.h b/services/source/transport/include/distributed_input_source_transport.h index c42c51a..4fbd9fe 100644 --- a/services/source/transport/include/distributed_input_source_transport.h +++ b/services/source/transport/include/distributed_input_source_transport.h @@ -22,6 +22,7 @@ #include #include #include +#include #include "constants.h" #include "event_handler.h" @@ -48,6 +49,10 @@ public: int32_t UnprepareRemoteInput(const std::string& deviceId); int32_t StartRemoteInput(const std::string& deviceId, const uint32_t& inputTypes); int32_t StopRemoteInput(const std::string& deviceId, const uint32_t& inputTypes); + int32_t LatencyCount(const std::string& deviceId); + void StartLatencyCount(const std::string& deviceId); + void StartLatencyThread(const std::string& deviceId); + void StopLatencyThread(); int32_t OnSessionOpened(int32_t sessionId, int32_t result); void OnSessionClosed(int32_t sessionId); @@ -58,11 +63,13 @@ private: int32_t SendMsg(int32_t sessionId, std::string &message); int32_t CheckDeviceSessionState(const std::string &remoteDevId); void HandleSessionData(int32_t sessionId, const std::string& messageData); + bool CheckRecivedData(const std::string& messageData); void NotifyResponsePrepareRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyResponseUnprepareRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyResponseStartRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyResponseStopRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); void NotifyReceivedEventRemoteInput(int32_t sessionId, const nlohmann::json &recMsg); + void CalculateLatency(int32_t sessionId, const nlohmann::json &recMsg); private: std::map sessionDevMap_; @@ -72,6 +79,14 @@ private: std::shared_ptr callback_; std::string mySessionName_ = ""; std::condition_variable openSessionWaitCond_; + uint64_t deltaTime_ = 0; + uint64_t deltaTimeAll_ = 0; + uint64_t sendTime_ = 0; + int32_t sendNum_ = 0; + int32_t recvNum_ = 0; + std::atomic isLatencyThreadRunning_ = false; + std::thread latencyThread_; + std::string eachLatencyDetails_ = ""; }; } // namespace DistributedInput } // namespace DistributedHardware diff --git a/services/source/transport/src/distributed_input_source_transport.cpp b/services/source/transport/src/distributed_input_source_transport.cpp index b7fe385..02ef269 100644 --- a/services/source/transport/src/distributed_input_source_transport.cpp +++ b/services/source/transport/src/distributed_input_source_transport.cpp @@ -29,6 +29,7 @@ #include "dinput_hitrace.h" #include "dinput_low_latency_utils.h" #include "dinput_softbus_define.h" +#include "dinput_utils_tool.h" #include "distributed_input_inject.h" #include "hidumper.h" #include "session.h" @@ -183,6 +184,8 @@ int32_t DistributedInputSourceTransport::OpenInputSoftbus(const std::string &rem } } + StartLatencyThread(remoteDevId); + DHLOGI("OpenSession success, remoteDevId:%s, sessionId:%s", GetAnonyString(remoteDevId).c_str(), GetAnonyInt32(sessionId).c_str()); DInputLowLatencyUtils::GetInstance().EnableSourceLowLatency(); @@ -201,6 +204,8 @@ void DistributedInputSourceTransport::CloseInputSoftbus(const std::string &remot } int32_t sessionId = sessionDevMap_[remoteDevId]; + StopLatencyThread(); + DHLOGI("RemoteDevId: %s, sessionId: %s", GetAnonyString(remoteDevId).c_str(), GetAnonyInt32(sessionId).c_str()); HiDumper::GetInstance().SetSessionStatus(remoteDevId, SessionStatus::CLOSING); CloseSession(sessionId); @@ -330,6 +335,71 @@ int32_t DistributedInputSourceTransport::StopRemoteInput( } } +int32_t DistributedInputSourceTransport::LatencyCount(const std::string& deviceId) +{ + std::unique_lock sessionLock(operationMutex_); + if (sessionDevMap_.count(deviceId) <= 0) { + DHLOGE("LatencyCount error, not find this device:%s.", GetAnonyString(deviceId).c_str()); + return ERR_DH_INPUT_SERVER_SOURCE_TRANSPORT_LATENCY_FAIL; + } + + int32_t sessionId = sessionDevMap_[deviceId]; + nlohmann::json jsonStr; + jsonStr[DINPUT_SOFTBUS_KEY_CMD_TYPE] = TRANS_SOURCE_MSG_LATENCY; + jsonStr[DINPUT_SOFTBUS_KEY_DEVICE_ID] = deviceId; + jsonStr[DINPUT_SOFTBUS_KEY_SESSION_ID] = sessionId; + std::string smsg = jsonStr.dump(); + int32_t ret = SendMsg(sessionId, smsg); + if (ret != DH_SUCCESS) { + DHLOGE("LatencyCount deviceId:%s, sessionId:%s, smsg:%s, SendMsg error, ret:%d.", + GetAnonyString(deviceId).c_str(), GetAnonyInt32(sessionId).c_str(), smsg.c_str(), ret); + return ERR_DH_INPUT_SERVER_SOURCE_TRANSPORT_LATENCY_FAIL; + } + + DHLOGI("LatencyCount deviceId:%s, sessionId:%s, smsg:%s.", + GetAnonyString(deviceId).c_str(), GetAnonyInt32(sessionId).c_str(), smsg.c_str()); + return DH_SUCCESS; +} + +void DistributedInputSourceTransport::StartLatencyCount(const std::string& deviceId) +{ + DHLOGI("start"); + while (isLatencyThreadRunning_.load()) { + if (sendNum_ >= INPUT_LATENCY_DELAY_TIMES) { + int32_t latency = deltaTimeAll_ / 2 / INPUT_LATENCY_DELAY_TIMES; + DHLOGI("LatencyCount average single-channel latency is %d, send times is %d, recive times is %d,\ + each RTT latency details is %s", latency, sendNum_, recvNum_, eachLatencyDetails_.c_str()); + deltaTimeAll_ = 0; + sendNum_ = 0; + recvNum_ = 0; + eachLatencyDetails_ = ""; + } + sendTime_ = GetCurrentTime(); + LatencyCount(deviceId); + sendNum_ += 1; + usleep(INPUT_LATENCY_DELAYTIME_US); + } + DHLOGI("end"); +} + +void DistributedInputSourceTransport::StartLatencyThread(const std::string& deviceId) +{ + DHLOGI("start"); + isLatencyThreadRunning_.store(true); + latencyThread_ = std::thread(&DistributedInputSourceTransport::StartLatencyCount, this, deviceId); + DHLOGI("end"); +} + +void DistributedInputSourceTransport::StopLatencyThread() +{ + DHLOGI("start"); + isLatencyThreadRunning_.store(false); + if (latencyThread_.joinable()) { + latencyThread_.join(); + } + DHLOGI("end"); +} + std::string DistributedInputSourceTransport::FindDeviceBySession(int32_t sessionId) { std::unique_lock sessionLock(operationMutex_); @@ -489,6 +559,21 @@ void DistributedInputSourceTransport::NotifyReceivedEventRemoteInput(int32_t ses callback_->onReceivedEventRemoteInput(deviceId, inputDataStr); } +void DistributedInputSourceTransport::CalculateLatency(int32_t sessionId, const nlohmann::json &recMsg) +{ + DHLOGI("OnBytesReceived cmdType is TRANS_SINK_MSG_LATENCY."); + std::string deviceId = FindDeviceBySession(sessionId); + if (deviceId.empty()) { + DHLOGE("OnBytesReceived cmdType is TRANS_SINK_MSG_LATENCY, deviceId is error."); + return; + } + + deltaTime_ = GetCurrentTime() - sendTime_; + deltaTimeAll_ += deltaTime_; + recvNum_ += 1; + eachLatencyDetails_ += (std::to_string(deltaTime_) + DINPUT_SPLIT_COMMA); +} + void DistributedInputSourceTransport::HandleSessionData(int32_t sessionId, const std::string& message) { if (callback_ == nullptr) { @@ -496,19 +581,7 @@ void DistributedInputSourceTransport::HandleSessionData(int32_t sessionId, const return; } nlohmann::json recMsg = nlohmann::json::parse(message); - if (recMsg.is_discarded()) { - DHLOGE("OnBytesReceived jsonStr error."); - return; - } - - if (recMsg.contains(DINPUT_SOFTBUS_KEY_CMD_TYPE) != true) { - DHLOGE("OnBytesReceived message:%s is error, not contain cmdType.", - message.c_str()); - return; - } - - if (recMsg[DINPUT_SOFTBUS_KEY_CMD_TYPE].is_number() != true) { - DHLOGE("OnBytesReceived cmdType is not number type."); + if (CheckRecivedData(message) != true) { return; } @@ -534,6 +607,10 @@ void DistributedInputSourceTransport::HandleSessionData(int32_t sessionId, const NotifyReceivedEventRemoteInput(sessionId, recMsg); break; } + case TRANS_SINK_MSG_LATENCY: { + CalculateLatency(sessionId, recMsg); + break; + } default: { DHLOGE("OnBytesReceived cmdType is undefined."); break; @@ -541,6 +618,28 @@ void DistributedInputSourceTransport::HandleSessionData(int32_t sessionId, const } } +bool DistributedInputSourceTransport::CheckRecivedData(const std::string& message) +{ + nlohmann::json recMsg = nlohmann::json::parse(message); + if (recMsg.is_discarded()) { + DHLOGE("OnBytesReceived jsonStr error."); + return false; + } + + if (recMsg.contains(DINPUT_SOFTBUS_KEY_CMD_TYPE) != true) { + DHLOGE("OnBytesReceived message:%s is error, not contain cmdType.", + message.c_str()); + return false; + } + + if (recMsg[DINPUT_SOFTBUS_KEY_CMD_TYPE].is_number() != true) { + DHLOGE("OnBytesReceived cmdType is not number type."); + return false; + } + + return true; +} + void DistributedInputSourceTransport::OnBytesReceived(int32_t sessionId, const void *data, uint32_t dataLen) { DHLOGI("OnBytesReceived, sessionId:%s, dataLen:%d", GetAnonyInt32(sessionId).c_str(), dataLen); diff --git a/utils/include/dinput_utils_tool.h b/utils/include/dinput_utils_tool.h index 070906e..71f25f0 100644 --- a/utils/include/dinput_utils_tool.h +++ b/utils/include/dinput_utils_tool.h @@ -31,6 +31,7 @@ struct DevInfo { }; DevInfo GetLocalDeviceInfo(); +uint64_t GetCurrentTime(); } } } diff --git a/utils/src/dinput_utils_tool.cpp b/utils/src/dinput_utils_tool.cpp index 8f33596..76f729d 100644 --- a/utils/src/dinput_utils_tool.cpp +++ b/utils/src/dinput_utils_tool.cpp @@ -15,6 +15,8 @@ #include "dinput_utils_tool.h" +#include + #include "softbus_bus_center.h" namespace OHOS { @@ -22,6 +24,7 @@ namespace DistributedHardware { namespace DistributedInput { namespace { const std::string DINPUT_PKG_NAME = "ohos.dhardware.dinput"; + constexpr int32_t MS_ONE_SECOND = 1000; } DevInfo GetLocalDeviceInfo() @@ -40,6 +43,13 @@ DevInfo GetLocalDeviceInfo() return devInfo; } + +uint64_t GetCurrentTime() +{ + struct timeval tv; + gettimeofday(&tv, nullptr); + return tv.tv_sec * MS_ONE_SECOND + tv.tv_usec / MS_ONE_SECOND; +} } } } \ No newline at end of file