mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-07-21 14:25:41 +00:00
refactor+feat: allow subsystems to send only declared messages, generate graphviz (#5314)
Closes #3774 Closes #3826
This commit is contained in:
committed by
GitHub
parent
26340b9054
commit
511891dcce
@@ -60,12 +60,13 @@
|
||||
#![deny(missing_docs)]
|
||||
#![deny(unused_crate_dependencies)]
|
||||
|
||||
pub use polkadot_overseer_gen_proc_macro::overlord;
|
||||
pub use polkadot_overseer_gen_proc_macro::{contextbounds, overlord, subsystem};
|
||||
|
||||
#[doc(hidden)]
|
||||
pub use gum;
|
||||
#[doc(hidden)]
|
||||
pub use metered;
|
||||
|
||||
#[doc(hidden)]
|
||||
pub use polkadot_node_primitives::SpawnNamed;
|
||||
|
||||
@@ -101,7 +102,7 @@ use std::fmt;
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
||||
/// A type of messages that are sent from [`Subsystem`] to [`Overseer`].
|
||||
/// A type of messages that are sent from a [`Subsystem`] to the declared overseer.
|
||||
///
|
||||
/// Used to launch jobs.
|
||||
pub enum ToOverseer {
|
||||
@@ -312,7 +313,7 @@ pub struct SubsystemMeterReadouts {
|
||||
///
|
||||
/// [`Subsystem`]: trait.Subsystem.html
|
||||
///
|
||||
/// `M` here is the inner message type, and _not_ the generated `enum AllMessages`.
|
||||
/// `M` here is the inner message type, and _not_ the generated `enum AllMessages` or `#message_wrapper` type.
|
||||
pub struct SubsystemInstance<Message, Signal> {
|
||||
/// Send sink for `Signal`s to be sent to a subsystem.
|
||||
pub tx_signal: crate::metered::MeteredSender<Signal>,
|
||||
@@ -362,20 +363,23 @@ pub trait SubsystemContext: Send + 'static {
|
||||
/// The message type of this context. Subsystems launched with this context will expect
|
||||
/// to receive messages of this type. Commonly uses the wrapping `enum` commonly called
|
||||
/// `AllMessages`.
|
||||
type Message: std::fmt::Debug + Send + 'static;
|
||||
type Message: ::std::fmt::Debug + Send + 'static;
|
||||
/// And the same for signals.
|
||||
type Signal: std::fmt::Debug + Send + 'static;
|
||||
/// The overarching all messages `enum`.
|
||||
/// In some cases can be identical to `Self::Message`.
|
||||
type AllMessages: From<Self::Message> + Send + 'static;
|
||||
type Signal: ::std::fmt::Debug + Send + 'static;
|
||||
/// The overarching messages `enum` for this particular subsystem.
|
||||
type OutgoingMessages: ::std::fmt::Debug + Send + 'static;
|
||||
|
||||
// The overarching messages `enum` for this particular subsystem.
|
||||
// type AllMessages: From<Self::OutgoingMessages> + From<Self::Message> + std::fmt::Debug + Send + 'static;
|
||||
|
||||
/// The sender type as provided by `sender()` and underlying.
|
||||
type Sender: SubsystemSender<Self::AllMessages> + Send + 'static;
|
||||
type Sender: Clone + Send + 'static + SubsystemSender<Self::OutgoingMessages>;
|
||||
/// The error type.
|
||||
type Error: ::std::error::Error + ::std::convert::From<OverseerError> + Sync + Send + 'static;
|
||||
|
||||
/// Try to asynchronously receive a message.
|
||||
///
|
||||
/// This has to be used with caution, if you loop over this without
|
||||
/// Has to be used with caution, if you loop over this without
|
||||
/// using `pending!()` macro you will end up with a busy loop!
|
||||
async fn try_recv(&mut self) -> Result<Option<FromOverseer<Self::Message, Self::Signal>>, ()>;
|
||||
|
||||
@@ -397,34 +401,37 @@ pub trait SubsystemContext: Send + 'static {
|
||||
) -> Result<(), Self::Error>;
|
||||
|
||||
/// Send a direct message to some other `Subsystem`, routed based on message type.
|
||||
async fn send_message<X>(&mut self, msg: X)
|
||||
// #[deprecated(note = "Use `self.sender().send_message(msg) instead, avoid passing around the full context.")]
|
||||
async fn send_message<T>(&mut self, msg: T)
|
||||
where
|
||||
Self::AllMessages: From<X>,
|
||||
X: Send,
|
||||
Self::OutgoingMessages: From<T> + Send,
|
||||
T: Send,
|
||||
{
|
||||
self.sender().send_message(<Self::AllMessages>::from(msg)).await
|
||||
self.sender().send_message(<Self::OutgoingMessages>::from(msg)).await
|
||||
}
|
||||
|
||||
/// Send multiple direct messages to other `Subsystem`s, routed based on message type.
|
||||
async fn send_messages<X, T>(&mut self, msgs: T)
|
||||
// #[deprecated(note = "Use `self.sender().send_message(msg) instead, avoid passing around the full context.")]
|
||||
async fn send_messages<T, I>(&mut self, msgs: I)
|
||||
where
|
||||
T: IntoIterator<Item = X> + Send,
|
||||
T::IntoIter: Send,
|
||||
Self::AllMessages: From<X>,
|
||||
X: Send,
|
||||
Self::OutgoingMessages: From<T> + Send,
|
||||
I: IntoIterator<Item = T> + Send,
|
||||
I::IntoIter: Send,
|
||||
T: Send,
|
||||
{
|
||||
self.sender()
|
||||
.send_messages(msgs.into_iter().map(|x| <Self::AllMessages>::from(x)))
|
||||
.send_messages(msgs.into_iter().map(<Self::OutgoingMessages>::from))
|
||||
.await
|
||||
}
|
||||
|
||||
/// Send a message using the unbounded connection.
|
||||
// #[deprecated(note = "Use `self.sender().send_unbounded_message(msg) instead, avoid passing around the full context.")]
|
||||
fn send_unbounded_message<X>(&mut self, msg: X)
|
||||
where
|
||||
Self::AllMessages: From<X>,
|
||||
Self::OutgoingMessages: From<X> + Send,
|
||||
X: Send,
|
||||
{
|
||||
self.sender().send_unbounded_message(Self::AllMessages::from(msg))
|
||||
self.sender().send_unbounded_message(<Self::OutgoingMessages>::from(msg))
|
||||
}
|
||||
|
||||
/// Obtain the sender.
|
||||
@@ -450,22 +457,25 @@ where
|
||||
|
||||
/// Sender end of a channel to interface with a subsystem.
|
||||
#[async_trait::async_trait]
|
||||
pub trait SubsystemSender<Message>: Send + Clone + 'static {
|
||||
pub trait SubsystemSender<OutgoingMessage>: Clone + Send + 'static
|
||||
where
|
||||
OutgoingMessage: Send,
|
||||
{
|
||||
/// Send a direct message to some other `Subsystem`, routed based on message type.
|
||||
async fn send_message(&mut self, msg: Message);
|
||||
async fn send_message(&mut self, msg: OutgoingMessage);
|
||||
|
||||
/// Send multiple direct messages to other `Subsystem`s, routed based on message type.
|
||||
async fn send_messages<T>(&mut self, msgs: T)
|
||||
async fn send_messages<I>(&mut self, msgs: I)
|
||||
where
|
||||
T: IntoIterator<Item = Message> + Send,
|
||||
T::IntoIter: Send;
|
||||
I: IntoIterator<Item = OutgoingMessage> + Send,
|
||||
I::IntoIter: Send;
|
||||
|
||||
/// Send a message onto the unbounded queue of some other `Subsystem`, routed based on message
|
||||
/// type.
|
||||
///
|
||||
/// This function should be used only when there is some other bounding factor on the messages
|
||||
/// sent with it. Otherwise, it risks a memory leak.
|
||||
fn send_unbounded_message(&mut self, msg: Message);
|
||||
fn send_unbounded_message(&mut self, msg: OutgoingMessage);
|
||||
}
|
||||
|
||||
/// A future that wraps another future with a `Delay` allowing for time-limited futures.
|
||||
|
||||
Reference in New Issue
Block a user