|
|
@@ -16,7 +16,7 @@
|
|
|
* along with this program. If not, see <https://www.gnu.org/licenses/>.
|
|
|
*/
|
|
|
|
|
|
-use std::sync::Arc;
|
|
|
+use std::{fmt, sync::Arc};
|
|
|
|
|
|
use darkfi_serial::{async_trait, serialize, SerialDecodable, SerialEncodable};
|
|
|
use log::{debug, error, info};
|
|
|
@@ -79,12 +79,6 @@ pub struct Channel {
|
|
|
pub info: ChannelInfo,
|
|
|
}
|
|
|
|
|
|
-impl std::fmt::Debug for Channel {
|
|
|
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
|
- write!(f, "{}", self.address())
|
|
|
- }
|
|
|
-}
|
|
|
-
|
|
|
impl Channel {
|
|
|
/// Sets up a new channel. Creates a reader and writer [`PtStream`] and
|
|
|
/// summons the message subscriber subsystem. Performs a network handshake
|
|
|
@@ -124,7 +118,7 @@ impl Channel {
|
|
|
/// Starts the channel. Runs a receive loop to start receiving messages
|
|
|
/// or handles a network failure.
|
|
|
pub fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
|
|
|
- debug!(target: "net::channel::start()", "START => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::start()", "START {:?}", self);
|
|
|
|
|
|
let self_ = self.clone();
|
|
|
self.receive_task.clone().start(
|
|
|
@@ -134,14 +128,14 @@ impl Channel {
|
|
|
executor,
|
|
|
);
|
|
|
|
|
|
- debug!(target: "net::channel::start()", "END => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::start()", "END {:?}", self);
|
|
|
}
|
|
|
|
|
|
/// Stops the channel. Steps through each component of the channel connection
|
|
|
/// and sends a stop signal. Notifies all subscribers that the channel has
|
|
|
/// been closed.
|
|
|
pub async fn stop(&self) {
|
|
|
- debug!(target: "net::channel::stop()", "START => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::stop()", "START {:?}", self);
|
|
|
|
|
|
if !*self.stopped.lock().await {
|
|
|
*self.stopped.lock().await = true;
|
|
|
@@ -151,13 +145,13 @@ impl Channel {
|
|
|
self.message_subsystem.trigger_error(Error::ChannelStopped).await;
|
|
|
}
|
|
|
|
|
|
- debug!(target: "net::channel::stop()", "END => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::stop()", "END {:?}", self);
|
|
|
}
|
|
|
|
|
|
/// Creates a subscription to a stopped signal.
|
|
|
/// If the channel is stopped then this will return a ChannelStopped error.
|
|
|
pub async fn subscribe_stop(&self) -> Result<Subscription<Error>> {
|
|
|
- debug!(target: "net::channel::subscribe_stop()", "START => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::subscribe_stop()", "START {:?}", self);
|
|
|
|
|
|
if *self.stopped.lock().await {
|
|
|
return Err(Error::ChannelStopped)
|
|
|
@@ -165,7 +159,7 @@ impl Channel {
|
|
|
|
|
|
let sub = self.stop_subscriber.clone().subscribe().await;
|
|
|
|
|
|
- debug!(target: "net::channel::subscribe_stop()", "END => address={}", self.address());
|
|
|
+ debug!(target: "net::channel::subscribe_stop()", "END {:?}", self);
|
|
|
|
|
|
Ok(sub)
|
|
|
}
|
|
|
@@ -175,7 +169,7 @@ impl Channel {
|
|
|
/// Returns an error if something goes wrong.
|
|
|
pub async fn send<M: message::Message>(&self, message: &M) -> Result<()> {
|
|
|
debug!(
|
|
|
- target: "net::channel::send()", "[START] command={} => address={}",
|
|
|
+ target: "net::channel::send()", "[START] command={} {:?}",
|
|
|
M::NAME, self.address(),
|
|
|
);
|
|
|
|
|
|
@@ -186,16 +180,16 @@ impl Channel {
|
|
|
// Catch failure and stop channel, return a net error
|
|
|
if let Err(e) = self.send_message(message).await {
|
|
|
error!(
|
|
|
- target: "net::channel::send()", "[P2P] Channel send error for [{}]: {}",
|
|
|
- self.address(), e
|
|
|
+ target: "net::channel::send()", "[P2P] Channel send error for [{:?}]: {}",
|
|
|
+ self, e
|
|
|
);
|
|
|
self.stop().await;
|
|
|
return Err(Error::ChannelStopped)
|
|
|
}
|
|
|
|
|
|
debug!(
|
|
|
- target: "net::channel::send()", "[END] command={} => address={}",
|
|
|
- M::NAME,self.address(),
|
|
|
+ target: "net::channel::send()", "[END] command={} {:?}",
|
|
|
+ M::NAME, self
|
|
|
);
|
|
|
|
|
|
Ok(())
|
|
|
@@ -223,15 +217,15 @@ impl Channel {
|
|
|
/// Subscribe to a message on the message subsystem.
|
|
|
pub async fn subscribe_msg<M: message::Message>(&self) -> Result<MessageSubscription<M>> {
|
|
|
debug!(
|
|
|
- target: "net::channel::subscribe_msg()", "[START] command={} => address={}",
|
|
|
- M::NAME, self.address(),
|
|
|
+ target: "net::channel::subscribe_msg()", "[START] command={} {:?}",
|
|
|
+ M::NAME, self
|
|
|
);
|
|
|
|
|
|
let sub = self.message_subsystem.subscribe::<M>().await;
|
|
|
|
|
|
debug!(
|
|
|
- target: "net::channel::subscribe_msg()", "[END] command={} => address={}",
|
|
|
- M::NAME, self.address(),
|
|
|
+ target: "net::channel::subscribe_msg()", "[END] command={} {:?}",
|
|
|
+ M::NAME, self
|
|
|
);
|
|
|
|
|
|
sub
|
|
|
@@ -240,7 +234,7 @@ impl Channel {
|
|
|
/// Handle network errors. Panic if error passes silently, otherwise
|
|
|
/// broadcast the error.
|
|
|
async fn handle_stop(self: Arc<Self>, result: Result<()>) {
|
|
|
- debug!(target: "net::channel::handle_stop()", "[START] address={}", self.address());
|
|
|
+ debug!(target: "net::channel::handle_stop()", "[START] {:?}", self);
|
|
|
|
|
|
match result {
|
|
|
Ok(()) => panic!("Channel task should never complete without error status"),
|
|
|
@@ -248,12 +242,12 @@ impl Channel {
|
|
|
Err(e) => self.message_subsystem.trigger_error(e).await,
|
|
|
}
|
|
|
|
|
|
- debug!(target: "net::channel::handle_stop()", "[END] address={}", self.address());
|
|
|
+ debug!(target: "net::channel::handle_stop()", "[END] {:?}", self);
|
|
|
}
|
|
|
|
|
|
/// Run the receive loop. Start receiving messages or handle network failure.
|
|
|
async fn main_receive_loop(self: Arc<Self>) -> Result<()> {
|
|
|
- debug!(target: "net::channel::main_receive_loop()", "[START] address={}", self.address());
|
|
|
+ debug!(target: "net::channel::main_receive_loop()", "[START] {:?}", self);
|
|
|
|
|
|
// Acquire reader lock
|
|
|
let reader = &mut *self.reader.lock().await;
|
|
|
@@ -279,7 +273,7 @@ impl Channel {
|
|
|
|
|
|
debug!(
|
|
|
target: "net::channel::main_receive_loop()",
|
|
|
- "Stopping channel {}", self.address(),
|
|
|
+ "Stopping channel {:?}", self
|
|
|
);
|
|
|
self.stop().await;
|
|
|
return Err(Error::ChannelStopped)
|
|
|
@@ -327,3 +321,9 @@ impl Channel {
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
+
|
|
|
+impl fmt::Debug for Channel {
|
|
|
+ fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
|
|
+ write!(f, "<Channel addr='{}' id={}>", self.info.addr, self.info.id)
|
|
|
+ }
|
|
|
+}
|