summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorHanjie Wu <hanjiew@andrew.cmu.edu>2023-12-18 05:30:17 -0500
committerHanjie Wu <hanjiew@andrew.cmu.edu>2023-12-18 05:30:17 -0500
commit47d92493a5bdd5a0bfb442b6238895a91d237c4c (patch)
treeadcfc92900b0c64c0888c187feba9eb98937c587
parentbd6528620f57fdc23cb4957f0879fb7301729198 (diff)
fix notification issue
-rw-r--r--src/PeerAbstractCommand.cc20
-rw-r--r--src/PeerAbstractCommand.h2
-rw-r--r--src/SocketCore.h2
-rw-r--r--src/uTPCore.cc34
-rw-r--r--src/uTPCore.h3
5 files changed, 42 insertions, 19 deletions
diff --git a/src/PeerAbstractCommand.cc b/src/PeerAbstractCommand.cc
index 03df3492..65995c98 100644
--- a/src/PeerAbstractCommand.cc
+++ b/src/PeerAbstractCommand.cc
@@ -126,7 +126,7 @@ bool PeerAbstractCommand::prepareForNextPeer(time_t wait) { return true; }
void PeerAbstractCommand::disableReadCheckSocket()
{
if (checkSocketIsReadable_) {
- e_->deleteSocketForReadCheck(readCheckTarget_, this);
+ setRWCheck(e_, readCheckTarget_, false, true);
checkSocketIsReadable_ = false;
readCheckTarget_.reset();
}
@@ -141,13 +141,13 @@ void PeerAbstractCommand::setReadCheckSocket(
else {
if (checkSocketIsReadable_) {
if (*readCheckTarget_ != *socket) {
- e_->deleteSocketForReadCheck(readCheckTarget_, this);
- e_->addSocketForReadCheck(socket, this);
+ setRWCheck(e_, readCheckTarget_, false, true);
+ setRWCheck(e_, socket, true, true);
readCheckTarget_ = socket;
}
}
else {
- e_->addSocketForReadCheck(socket, this);
+ setRWCheck(e_, socket, true, true);
checkSocketIsReadable_ = true;
readCheckTarget_ = socket;
}
@@ -157,7 +157,7 @@ void PeerAbstractCommand::setReadCheckSocket(
void PeerAbstractCommand::disableWriteCheckSocket()
{
if (checkSocketIsWritable_) {
- e_->deleteSocketForWriteCheck(writeCheckTarget_, this);
+ setRWCheck(e_, writeCheckTarget_, false, false);
checkSocketIsWritable_ = false;
writeCheckTarget_.reset();
}
@@ -172,20 +172,20 @@ void PeerAbstractCommand::setWriteCheckSocket(
else {
if (checkSocketIsWritable_) {
if (*writeCheckTarget_ != *socket) {
- e_->deleteSocketForWriteCheck(writeCheckTarget_, this);
- e_->addSocketForWriteCheck(socket, this);
+ setRWCheck(e_, writeCheckTarget_, false, false);
+ setRWCheck(e_, socket, true, false);
writeCheckTarget_ = socket;
}
}
else {
- e_->addSocketForWriteCheck(socket, this);
+ setRWCheck(e_, socket, true, false);
checkSocketIsWritable_ = true;
writeCheckTarget_ = socket;
}
}
}
-void PeerAbstractCommand::SetRWCheck(DownloadEngine* e,
+void PeerAbstractCommand::setRWCheck(DownloadEngine* e,
const std::shared_ptr<SocketCore>& socket,
int is_add, int is_read)
{
@@ -201,7 +201,7 @@ void PeerAbstractCommand::SetRWCheck(DownloadEngine* e,
}
}
else {
- utp->SetRWCheck(this, is_add, is_read);
+ utp->setRWCheck(this, is_add, is_read);
}
}
diff --git a/src/PeerAbstractCommand.h b/src/PeerAbstractCommand.h
index acc305a1..945cc39b 100644
--- a/src/PeerAbstractCommand.h
+++ b/src/PeerAbstractCommand.h
@@ -88,7 +88,7 @@ protected:
void setWriteCheckSocket(const std::shared_ptr<SocketCore>& socket);
void disableReadCheckSocket();
void disableWriteCheckSocket();
- void SetRWCheck(DownloadEngine* e, const std::shared_ptr<SocketCore>& socket,
+ void setRWCheck(DownloadEngine* e, const std::shared_ptr<SocketCore>& socket,
int is_add, int is_read);
void setNoCheck(bool check);
void updateKeepAlive();
diff --git a/src/SocketCore.h b/src/SocketCore.h
index b2dda052..3b28df4f 100644
--- a/src/SocketCore.h
+++ b/src/SocketCore.h
@@ -158,7 +158,7 @@ public:
sock_t getSockfd() const { return sockfd_; }
- bool isOpen() const { return sockfd_ != (sock_t)-1; }
+ virtual bool isOpen() const { return sockfd_ != (sock_t)-1; }
void setMulticastInterface(const std::string& localAddr);
diff --git a/src/uTPCore.cc b/src/uTPCore.cc
index a279ae56..a983676c 100644
--- a/src/uTPCore.cc
+++ b/src/uTPCore.cc
@@ -301,6 +301,34 @@ void uTPCore::closeConnection()
command_ = nullptr;
}
+void uTPCore::setRWCheck(Command* command, int is_add, int is_read)
+{
+ if (is_add == true) {
+ command_ = command;
+ if (is_read == true) {
+ readCheck_ = true;
+ if (isReadable(0)) {
+ command_->setStatusActive();
+ }
+ }
+ else {
+ writeCheck_ = true;
+ if (isWritable(0)) {
+ command_->setStatusActive();
+ }
+ }
+ }
+ else {
+ command_ = command;
+ if (is_read == true) {
+ readCheck_ = false;
+ }
+ else {
+ writeCheck_ = false;
+ }
+ }
+}
+
uTPListenCommand::uTPListenCommand(cuid_t cuid, DownloadEngine* e,
utp_context* ctx,
std::shared_ptr<SocketCore> sock)
@@ -344,10 +372,4 @@ bool uTPListenCommand::execute()
return false;
}
-void uTPCore::SetRWCheck(Command* command, int is_add, int is_read)
-{
- command_ = command;
- (is_read ? readCheck_ : writeCheck_) = is_add;
-}
-
} // namespace aria2
diff --git a/src/uTPCore.h b/src/uTPCore.h
index afc13fd6..44a73aa5 100644
--- a/src/uTPCore.h
+++ b/src/uTPCore.h
@@ -44,6 +44,7 @@ public:
virtual void establishConnection(const std::string& host, uint16_t port,
bool tcpNodelay = true) CXX11_OVERRIDE;
+ virtual bool isOpen() const CXX11_OVERRIDE { return s_ != nullptr; };
virtual bool isWritable(time_t timeout) CXX11_OVERRIDE;
virtual bool isReadable(time_t timeout) CXX11_OVERRIDE;
virtual ssize_t writeData(const void* data, size_t len) CXX11_OVERRIDE;
@@ -51,7 +52,7 @@ public:
virtual void readData(void* data, size_t& len) CXX11_OVERRIDE;
virtual void closeConnection() CXX11_OVERRIDE;
- void SetRWCheck(Command* command, int is_add, int is_read);
+ void setRWCheck(Command* command, int is_add, int is_read);
static uint64 getMilliseconds(utp_callback_arguments* arg);
static uint64 getMicroseconds(utp_callback_arguments* arg);