diff options
| author | Hanjie Wu <hanjiew@andrew.cmu.edu> | 2023-12-18 07:52:27 -0500 |
|---|---|---|
| committer | Hanjie Wu <hanjiew@andrew.cmu.edu> | 2023-12-18 07:52:27 -0500 |
| commit | 2d3d731726d77099fe50a0e3c228af2263d42ffb (patch) | |
| tree | 96f9277e919db59ca5256f45984196859fe7e450 | |
| parent | 47d92493a5bdd5a0bfb442b6238895a91d237c4c (diff) | |
fix notification
| -rw-r--r-- | src/SocketCore.h | 2 | ||||
| -rw-r--r-- | src/uTPCore.cc | 29 | ||||
| -rw-r--r-- | src/uTPCore.h | 1 |
3 files changed, 26 insertions, 6 deletions
diff --git a/src/SocketCore.h b/src/SocketCore.h index 3b28df4f..14b26a6f 100644 --- a/src/SocketCore.h +++ b/src/SocketCore.h @@ -369,7 +369,7 @@ public: bool operator<(const SocketCore& s) { return sockfd_ < s.sockfd_; } - std::string getSocketError() const; + virtual std::string getSocketError() const; /** * Returns true if the underlying socket gets EAGAIN in the previous diff --git a/src/uTPCore.cc b/src/uTPCore.cc index a983676c..4113f959 100644 --- a/src/uTPCore.cc +++ b/src/uTPCore.cc @@ -98,6 +98,7 @@ uint64 uTPCore::stateChangeCallback(utp_callback_arguments* arg) utp->isWritable_ = true; if (utp->command_ != nullptr && utp->writeCheck_ == true) { utp->command_->setStatusActive(); + utp->command_->writeEventReceived(); } } return 0; @@ -120,6 +121,7 @@ uint64 uTPCore::readCallback(utp_callback_arguments* arg) } if (utp->command_ != nullptr && utp->readCheck_ == true) { utp->command_->setStatusActive(); + utp->command_->readEventReceived(); } return 0; } @@ -127,9 +129,12 @@ uint64 uTPCore::readCallback(utp_callback_arguments* arg) uint64 uTPCore::errorCallback(utp_callback_arguments* arg) { uTPCore* utp = (uTPCore*)utp_get_userdata(arg->socket); - utp->errNum_ = arg->error_code; - if (utp->command_ != nullptr) { - utp->command_->setStatusActive(); + if (utp != nullptr) { + utp->errNum_ = arg->error_code + 1; + if (utp->command_ != nullptr) { + utp->command_->setStatusActive(); + utp->command_->errorEventReceived(); + } } return 0; } @@ -156,6 +161,7 @@ void setupContextAndSock(utp_context*& ctx, std::shared_ptr<SocketCore>& sock, utp_set_callback(ctx, UTP_ON_ACCEPT, uTPCore::acceptCallback); utp_set_callback(ctx, UTP_ON_STATE_CHANGE, uTPCore::stateChangeCallback); utp_set_callback(ctx, UTP_ON_READ, uTPCore::readCallback); + utp_set_callback(ctx, UTP_ON_ERROR, uTPCore::errorCallback); e->addRoutineCommand( make_unique<uTPListenCommand>(e->newCUID(), e, ctx, sock)); @@ -241,7 +247,7 @@ ssize_t uTPCore::writeVector(a2iovec* iov, size_t iovcnt) { if (errNum_ != 0) { throw DL_RETRY_EX( - fmt("Failed to send data, cause: %s.", utp_error_code_names[errNum_])); + fmt("Failed to send data, cause: %s.", getSocketError().c_str())); } ssize_t ret = 0; wantRead_ = false; @@ -264,7 +270,7 @@ void uTPCore::readData(void* data, size_t& len) { if (errNum_ != 0) { throw DL_RETRY_EX( - fmt("Failed to read data, cause: %s.", utp_error_code_names[errNum_])); + fmt("Failed to read data, cause: %s.", getSocketError().c_str())); } wantRead_ = false; wantWrite_ = false; @@ -301,6 +307,17 @@ void uTPCore::closeConnection() command_ = nullptr; } +std::string uTPCore::getSocketError() const +{ + if (errNum_ == 0) { + return "NO_ERROR"; + } + else if (errNum_ >= 1 && errNum_ <= 3) { + return utp_error_code_names[errNum_ - 1]; + } + return "UNKNOWN_ERROR"; +}; + void uTPCore::setRWCheck(Command* command, int is_add, int is_read) { if (is_add == true) { @@ -309,12 +326,14 @@ void uTPCore::setRWCheck(Command* command, int is_add, int is_read) readCheck_ = true; if (isReadable(0)) { command_->setStatusActive(); + command_->readEventReceived(); } } else { writeCheck_ = true; if (isWritable(0)) { command_->setStatusActive(); + command_->writeEventReceived(); } } } diff --git a/src/uTPCore.h b/src/uTPCore.h index 44a73aa5..838d0500 100644 --- a/src/uTPCore.h +++ b/src/uTPCore.h @@ -51,6 +51,7 @@ public: virtual ssize_t writeVector(a2iovec* iov, size_t iovcnt) CXX11_OVERRIDE; virtual void readData(void* data, size_t& len) CXX11_OVERRIDE; virtual void closeConnection() CXX11_OVERRIDE; + virtual std::string getSocketError() const CXX11_OVERRIDE; void setRWCheck(Command* command, int is_add, int is_read); |
