chore: clean up emitter-and-signal library design
This commit is contained in:
@@ -1,9 +1,9 @@
|
||||
use snafu::Snafu;
|
||||
|
||||
pub mod emitter;
|
||||
pub mod emitter_ext;
|
||||
mod emitter_ext;
|
||||
pub mod signal;
|
||||
pub mod signal_ext;
|
||||
mod signal_ext;
|
||||
|
||||
pub use emitter::Emitter;
|
||||
pub use emitter_ext::EmitterExt;
|
||||
|
||||
@@ -1,30 +1,24 @@
|
||||
use std::future::Future;
|
||||
|
||||
use ext_trait::extension;
|
||||
use snafu::{ResultExt, Snafu};
|
||||
use tokio::select;
|
||||
|
||||
use crate::ProducerExited;
|
||||
|
||||
use super::signal::{JoinError, Signal};
|
||||
|
||||
#[derive(Debug, Snafu)]
|
||||
pub struct ProducerAlreadyExited {
|
||||
source: ProducerExited,
|
||||
}
|
||||
|
||||
#[extension(pub trait SignalExt)]
|
||||
impl<T> Signal<T> {
|
||||
fn map<M, F>(
|
||||
self,
|
||||
mut func: F,
|
||||
) -> Result<(Signal<M>, impl Future<Output = Result<(), JoinError>>), ProducerAlreadyExited>
|
||||
) -> Result<(Signal<M>, impl Future<Output = Result<(), JoinError>>), ProducerExited>
|
||||
where
|
||||
T: 'static + Sync + Send + Clone,
|
||||
M: 'static + Sync + Send + Clone,
|
||||
F: 'static + Send + FnMut(T) -> M,
|
||||
{
|
||||
let initial = func(self.subscribe().context(ProducerAlreadyExitedSnafu)?.get());
|
||||
let initial = func(self.subscribe()?.get());
|
||||
|
||||
Ok(Signal::new(initial, |mut publisher_stream| async move {
|
||||
while let Some(publisher) = publisher_stream.wait().await {
|
||||
|
||||
Reference in New Issue
Block a user