mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-07-29 06:45:42 +00:00
This reverts commit 6c9c687a31.
Co-authored-by: André Silva <andrerfosilva@gmail.com>
This commit is contained in:
@@ -88,11 +88,10 @@ enum NotificationsInSubstreamHandshake {
|
||||
PendingSend(Vec<u8>),
|
||||
/// Handshake message was pushed in the socket. Still need to flush.
|
||||
Flush,
|
||||
/// Ready to receive notifications. Handshake message successfully sent and flushed, or
|
||||
/// sending side closed before handshake sent.
|
||||
Normal { write_side_open: bool },
|
||||
/// Closing our writing side.
|
||||
Closing { remote_write_open: bool },
|
||||
/// Handshake message successfully sent and flushed.
|
||||
Sent,
|
||||
/// Remote has closed their writing side. We close our own writing side in return.
|
||||
ClosingInResponseToRemote,
|
||||
/// Both our side and the remote have closed their writing side.
|
||||
BothSidesClosed,
|
||||
}
|
||||
@@ -170,30 +169,8 @@ where TSubstream: AsyncRead + AsyncWrite + Unpin + Send + 'static,
|
||||
impl<TSubstream> NotificationsInSubstream<TSubstream>
|
||||
where TSubstream: AsyncRead + AsyncWrite + Unpin,
|
||||
{
|
||||
/// Closes the writing side of the substream, indicating to the remote that we would like this
|
||||
/// substream to be closed.
|
||||
pub fn set_close_desired(&mut self) {
|
||||
match self.handshake {
|
||||
NotificationsInSubstreamHandshake::PendingSend(_) |
|
||||
NotificationsInSubstreamHandshake::Flush |
|
||||
NotificationsInSubstreamHandshake::NotSent |
|
||||
NotificationsInSubstreamHandshake::Normal { write_side_open: true } => {
|
||||
self.handshake = NotificationsInSubstreamHandshake::Closing { remote_write_open: true };
|
||||
}
|
||||
NotificationsInSubstreamHandshake::Normal { write_side_open: false } |
|
||||
NotificationsInSubstreamHandshake::Closing { .. } |
|
||||
NotificationsInSubstreamHandshake::BothSidesClosed => {}
|
||||
}
|
||||
}
|
||||
|
||||
/// Sends the handshake in order to inform the remote that we accept the substream.
|
||||
///
|
||||
/// Has no effect if `set_close_desired` has been called.
|
||||
pub fn send_handshake(&mut self, message: impl Into<Vec<u8>>) {
|
||||
if matches!(self.handshake, NotificationsInSubstreamHandshake::Normal { write_side_open: false }) {
|
||||
return;
|
||||
}
|
||||
|
||||
if !matches!(self.handshake, NotificationsInSubstreamHandshake::NotSent) {
|
||||
error!(target: "sub-libp2p", "Tried to send handshake twice");
|
||||
return;
|
||||
@@ -208,7 +185,7 @@ where TSubstream: AsyncRead + AsyncWrite + Unpin,
|
||||
let mut this = self.project();
|
||||
|
||||
loop {
|
||||
match mem::replace(this.handshake, NotificationsInSubstreamHandshake::NotSent) {
|
||||
match mem::replace(this.handshake, NotificationsInSubstreamHandshake::Sent) {
|
||||
NotificationsInSubstreamHandshake::PendingSend(msg) =>
|
||||
match Sink::poll_ready(this.socket.as_mut(), cx) {
|
||||
Poll::Ready(_) => {
|
||||
@@ -226,28 +203,16 @@ where TSubstream: AsyncRead + AsyncWrite + Unpin,
|
||||
NotificationsInSubstreamHandshake::Flush =>
|
||||
match Sink::poll_flush(this.socket.as_mut(), cx)? {
|
||||
Poll::Ready(()) =>
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open: true },
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Sent,
|
||||
Poll::Pending => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Flush;
|
||||
return Poll::Pending
|
||||
}
|
||||
},
|
||||
|
||||
NotificationsInSubstreamHandshake::Closing { remote_write_open } =>
|
||||
match Sink::poll_close(this.socket.as_mut(), cx)? {
|
||||
Poll::Ready(()) => if remote_write_open {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open: false }
|
||||
} else {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::BothSidesClosed;
|
||||
},
|
||||
Poll::Pending => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Closing { remote_write_open };
|
||||
return Poll::Pending
|
||||
}
|
||||
},
|
||||
|
||||
st @ NotificationsInSubstreamHandshake::NotSent |
|
||||
st @ NotificationsInSubstreamHandshake::Normal { .. } |
|
||||
st @ NotificationsInSubstreamHandshake::Sent |
|
||||
st @ NotificationsInSubstreamHandshake::ClosingInResponseToRemote |
|
||||
st @ NotificationsInSubstreamHandshake::BothSidesClosed => {
|
||||
*this.handshake = st;
|
||||
return Poll::Pending;
|
||||
@@ -267,7 +232,7 @@ where TSubstream: AsyncRead + AsyncWrite + Unpin,
|
||||
|
||||
// This `Stream` implementation first tries to send back the handshake if necessary.
|
||||
loop {
|
||||
match mem::replace(this.handshake, NotificationsInSubstreamHandshake::NotSent) {
|
||||
match mem::replace(this.handshake, NotificationsInSubstreamHandshake::Sent) {
|
||||
NotificationsInSubstreamHandshake::NotSent => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::NotSent;
|
||||
return Poll::Pending
|
||||
@@ -289,39 +254,34 @@ where TSubstream: AsyncRead + AsyncWrite + Unpin,
|
||||
NotificationsInSubstreamHandshake::Flush =>
|
||||
match Sink::poll_flush(this.socket.as_mut(), cx)? {
|
||||
Poll::Ready(()) =>
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open: true },
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Sent,
|
||||
Poll::Pending => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Flush;
|
||||
return Poll::Pending
|
||||
}
|
||||
},
|
||||
|
||||
NotificationsInSubstreamHandshake::Normal { write_side_open } => {
|
||||
NotificationsInSubstreamHandshake::Sent => {
|
||||
match Stream::poll_next(this.socket.as_mut(), cx) {
|
||||
Poll::Ready(None) if write_side_open =>
|
||||
*this.handshake =
|
||||
NotificationsInSubstreamHandshake::Closing { remote_write_open: false },
|
||||
Poll::Ready(None) =>
|
||||
*this.handshake = NotificationsInSubstreamHandshake::BothSidesClosed,
|
||||
Poll::Ready(None) => *this.handshake =
|
||||
NotificationsInSubstreamHandshake::ClosingInResponseToRemote,
|
||||
Poll::Ready(Some(msg)) => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open };
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Sent;
|
||||
return Poll::Ready(Some(msg))
|
||||
},
|
||||
Poll::Pending => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open };
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Sent;
|
||||
return Poll::Pending
|
||||
},
|
||||
}
|
||||
},
|
||||
|
||||
NotificationsInSubstreamHandshake::Closing { remote_write_open } =>
|
||||
NotificationsInSubstreamHandshake::ClosingInResponseToRemote =>
|
||||
match Sink::poll_close(this.socket.as_mut(), cx)? {
|
||||
Poll::Ready(()) if remote_write_open =>
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Normal { write_side_open: false },
|
||||
Poll::Ready(()) =>
|
||||
*this.handshake = NotificationsInSubstreamHandshake::BothSidesClosed,
|
||||
Poll::Pending => {
|
||||
*this.handshake = NotificationsInSubstreamHandshake::Closing { remote_write_open };
|
||||
*this.handshake = NotificationsInSubstreamHandshake::ClosingInResponseToRemote;
|
||||
return Poll::Pending
|
||||
}
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user