v0.18.21: format code
This commit is contained in:
+4
-1
@@ -27,4 +27,7 @@ pub(crate) use writeback_file::WritebackFile;
|
|||||||
// caller — the targeted `#[allow]` keeps the re-export visible without
|
// caller — the targeted `#[allow]` keeps the re-export visible without
|
||||||
// dragging the rest of the module under `dead_code`.
|
// dragging the rest of the module under `dead_code`.
|
||||||
#[allow(unused_imports)]
|
#[allow(unused_imports)]
|
||||||
pub use pipeline::{DEFAULT_PIPELINE_DEPTH, Flow, Pipeline, READ_PIPELINE_DEPTH, Sink, WRITE_PIPELINE_DEPTH, WRITE_THROUGH_DEPTH};
|
pub use pipeline::{
|
||||||
|
DEFAULT_PIPELINE_DEPTH, Flow, Pipeline, READ_PIPELINE_DEPTH, Sink, WRITE_PIPELINE_DEPTH,
|
||||||
|
WRITE_THROUGH_DEPTH,
|
||||||
|
};
|
||||||
|
|||||||
+38
-44
@@ -153,57 +153,51 @@ impl<I: Send + 'static, R: Send + 'static> Pipeline<I, R> {
|
|||||||
let mut first_err: Option<Error> = None;
|
let mut first_err: Option<Error> = None;
|
||||||
let mut stopped = false;
|
let mut stopped = false;
|
||||||
|
|
||||||
while let Ok(item) = rx.recv() {
|
while let Ok(item) = rx.recv() {
|
||||||
if debug_enabled() {
|
if debug_enabled() {
|
||||||
tracing::debug!(
|
tracing::debug!("Pipeline receive: item={}", std::any::type_name::<I>());
|
||||||
"Pipeline receive: item={}",
|
}
|
||||||
std::any::type_name::<I>()
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
let apply_start = std::time::Instant::now();
|
let apply_start = std::time::Instant::now();
|
||||||
|
|
||||||
if first_err.is_some() || stopped {
|
if first_err.is_some() || stopped {
|
||||||
// Drain remaining items so the producer never
|
// Drain remaining items so the producer never
|
||||||
// blocks on a dead receiver. `apply` is not
|
// blocks on a dead receiver. `apply` is not
|
||||||
// called once we've decided to stop.
|
// called once we've decided to stop.
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
match sink.apply(item) {
|
match sink.apply(item) {
|
||||||
Ok(Flow::Continue) => {}
|
Ok(Flow::Continue) => {}
|
||||||
Ok(Flow::Stop) => {
|
Ok(Flow::Stop) => {
|
||||||
stopped = true;
|
stopped = true;
|
||||||
if debug_enabled() {
|
if debug_enabled() {
|
||||||
tracing::debug!("Pipeline: consumer returned Flow::Stop");
|
tracing::debug!("Pipeline: consumer returned Flow::Stop");
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
if debug_enabled() {
|
|
||||||
tracing::debug!(
|
|
||||||
"Pipeline: apply error, stopping, err={:?}",
|
|
||||||
e
|
|
||||||
);
|
|
||||||
}
|
|
||||||
first_err = Some(e);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
Err(e) => {
|
||||||
let apply_elapsed = apply_start.elapsed();
|
if debug_enabled() {
|
||||||
if debug_enabled() && apply_elapsed > std::time::Duration::from_millis(100) {
|
tracing::debug!("Pipeline: apply error, stopping, err={:?}", e);
|
||||||
tracing::debug!(
|
}
|
||||||
"Pipeline apply: took {:.2}s, item={}",
|
first_err = Some(e);
|
||||||
apply_elapsed.as_secs_f64(),
|
|
||||||
std::any::type_name::<I>()
|
|
||||||
);
|
|
||||||
} else if debug_enabled() {
|
|
||||||
tracing::debug!(
|
|
||||||
"Pipeline apply: OK in {:.3}ms, item={}",
|
|
||||||
apply_elapsed.as_micros(),
|
|
||||||
std::any::type_name::<I>()
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let apply_elapsed = apply_start.elapsed();
|
||||||
|
if debug_enabled() && apply_elapsed > std::time::Duration::from_millis(100) {
|
||||||
|
tracing::debug!(
|
||||||
|
"Pipeline apply: took {:.2}s, item={}",
|
||||||
|
apply_elapsed.as_secs_f64(),
|
||||||
|
std::any::type_name::<I>()
|
||||||
|
);
|
||||||
|
} else if debug_enabled() {
|
||||||
|
tracing::debug!(
|
||||||
|
"Pipeline apply: OK in {:.3}ms, item={}",
|
||||||
|
apply_elapsed.as_micros(),
|
||||||
|
std::any::type_name::<I>()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
match first_err {
|
match first_err {
|
||||||
Some(e) => Err(e),
|
Some(e) => Err(e),
|
||||||
None => sink.close(),
|
None => sink.close(),
|
||||||
|
|||||||
+12
-12
@@ -72,18 +72,18 @@
|
|||||||
//! | E7xxx | AACS errors |
|
//! | E7xxx | AACS errors |
|
||||||
|
|
||||||
pub mod aacs;
|
pub mod aacs;
|
||||||
pub(crate) mod clpi;
|
pub(crate) mod clpi;
|
||||||
pub mod css;
|
pub mod css;
|
||||||
pub mod decrypt;
|
pub mod decrypt;
|
||||||
pub mod disc;
|
pub mod disc;
|
||||||
pub mod drive;
|
pub mod drive;
|
||||||
pub mod error;
|
pub mod error;
|
||||||
pub mod event;
|
pub mod event;
|
||||||
pub mod halt;
|
pub mod halt;
|
||||||
pub(crate) mod identity;
|
pub(crate) mod identity;
|
||||||
pub(crate) mod ifo;
|
pub(crate) mod ifo;
|
||||||
pub mod io;
|
pub mod io;
|
||||||
pub mod keydb;
|
pub mod keydb;
|
||||||
pub mod labels;
|
pub mod labels;
|
||||||
pub(crate) mod mpls;
|
pub(crate) mod mpls;
|
||||||
pub mod mux;
|
pub mod mux;
|
||||||
|
|||||||
Reference in New Issue
Block a user