summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorHanjie Wu <hanjiew@andrew.cmu.edu>2023-12-18 07:52:27 -0500
committerHanjie Wu <hanjiew@andrew.cmu.edu>2023-12-18 07:52:27 -0500
commit2d3d731726d77099fe50a0e3c228af2263d42ffb (patch)
tree96f9277e919db59ca5256f45984196859fe7e450
parent47d92493a5bdd5a0bfb442b6238895a91d237c4c (diff)
fix notification
-rw-r--r--src/SocketCore.h2
-rw-r--r--src/uTPCore.cc29
-rw-r--r--src/uTPCore.h1
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);