This commit is contained in:
Dominik Werder
2022-11-29 16:37:09 +01:00
parent 9fe63706cf
commit 047237e250
20 changed files with 266 additions and 176 deletions

View File

@@ -41,6 +41,7 @@ bitshuffle = { path = "../bitshuffle" }
dbconn = { path = "../dbconn" }
parse = { path = "../parse" }
items = { path = "../items" }
items_0 = { path = "../items_0" }
items_2 = { path = "../items_2" }
streams = { path = "../streams" }
httpclient = { path = "../httpclient" }

View File

@@ -11,6 +11,7 @@ use netpod::log::*;
use netpod::query::RawEventsQuery;
use netpod::{AggKind, ByteOrder, ByteSize, Channel, DiskIoTune, NanoRange, NodeConfigCached, ScalarType, Shape};
use parse::channelconfig::{extract_matching_config_entry, read_local_config, ConfigEntry, MatchingConfigEntry};
use std::collections::VecDeque;
use std::pin::Pin;
use streams::eventchunker::EventChunkerConf;
@@ -33,7 +34,16 @@ where
StreamItem::DataItem(item) => match item {
RangeCompletableItem::Data(item) => {
let item = events_node_proc.process(item);
Ok(StreamItem::DataItem(RangeCompletableItem::Data(item)))
use items::EventsNodeProcessorOutput;
let parts = item.into_parts::<NTY>();
let item = items_2::eventsdim0::EventsDim0 {
tss: parts.1,
pulses: VecDeque::new(),
values: parts.0,
};
let item = Box::new(item) as Box<dyn items_0::Events>;
//Ok(StreamItem::DataItem(RangeCompletableItem::Data(todo!())))
Ok(StreamItem::DataItem(RangeCompletableItem::RangeComplete))
}
RangeCompletableItem::RangeComplete => Ok(StreamItem::DataItem(RangeCompletableItem::RangeComplete)),
},