This commit is contained in:
Dominik Werder
2023-04-18 13:18:54 +02:00
parent b3f53c60d8
commit abd9cfdf65
+3 -5
View File
@@ -46,12 +46,11 @@ macro_rules! trace4 {
pub struct Collect { pub struct Collect {
inp: CollectableStreamBox, inp: CollectableStreamBox,
deadline: Instant,
events_max: u64, events_max: u64,
range: Option<SeriesRange>, range: Option<SeriesRange>,
binrange: Option<BinnedRangeEnum>, binrange: Option<BinnedRangeEnum>,
collector: Option<Box<dyn Collector>>, collector: Option<Box<dyn Collector>>,
range_complete: bool, range_final: bool,
timeout: bool, timeout: bool,
timer: Pin<Box<dyn Future<Output = ()> + Send>>, timer: Pin<Box<dyn Future<Output = ()> + Send>>,
done_input: bool, done_input: bool,
@@ -71,12 +70,11 @@ impl Collect {
let timer = tokio::time::sleep_until(deadline.into()); let timer = tokio::time::sleep_until(deadline.into());
Self { Self {
inp: CollectableStreamBox(Box::pin(inp)), inp: CollectableStreamBox(Box::pin(inp)),
deadline,
events_max, events_max,
range, range,
binrange, binrange,
collector: None, collector: None,
range_complete: false, range_final: false,
timeout: false, timeout: false,
timer: Box::pin(timer), timer: Box::pin(timer),
done_input: false, done_input: false,
@@ -88,7 +86,7 @@ impl Collect {
Ok(item) => match item { Ok(item) => match item {
StreamItem::DataItem(item) => match item { StreamItem::DataItem(item) => match item {
RangeCompletableItem::RangeComplete => { RangeCompletableItem::RangeComplete => {
self.range_complete = true; self.range_final = true;
if let Some(coll) = self.collector.as_mut() { if let Some(coll) = self.collector.as_mut() {
coll.set_range_complete(); coll.set_range_complete();
} else { } else {