diff options
| author | Hanjie Wu <hanjiew@andrew.cmu.edu> | 2023-12-18 05:30:17 -0500 |
|---|---|---|
| committer | Hanjie Wu <hanjiew@andrew.cmu.edu> | 2023-12-18 05:30:17 -0500 |
| commit | 47d92493a5bdd5a0bfb442b6238895a91d237c4c (patch) | |
| tree | adcfc92900b0c64c0888c187feba9eb98937c587 | |
| parent | bd6528620f57fdc23cb4957f0879fb7301729198 (diff) | |
fix notification issue
| -rw-r--r-- | src/PeerAbstractCommand.cc | 20 | ||||
| -rw-r--r-- | src/PeerAbstractCommand.h | 2 | ||||
| -rw-r--r-- | src/SocketCore.h | 2 | ||||
| -rw-r--r-- | src/uTPCore.cc | 34 | ||||
| -rw-r--r-- | src/uTPCore.h | 3 |
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); |
