Check for requested backend
This commit is contained in:
@@ -18,6 +18,7 @@ pub fn make_test_node(id: u32) -> Node {
|
|||||||
data_base_path: format!("../tmpdata/node{:02}", id).into(),
|
data_base_path: format!("../tmpdata/node{:02}", id).into(),
|
||||||
split: id,
|
split: id,
|
||||||
ksprefix: "ks".into(),
|
ksprefix: "ks".into(),
|
||||||
|
backend: "testbackend".into(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -158,6 +158,13 @@ pub async fn binned_bytes_for_http(
|
|||||||
node_config: &NodeConfigCached,
|
node_config: &NodeConfigCached,
|
||||||
query: &BinnedQuery,
|
query: &BinnedQuery,
|
||||||
) -> Result<BinnedBytesForHttpStream, Error> {
|
) -> Result<BinnedBytesForHttpStream, Error> {
|
||||||
|
if query.channel.backend != node_config.node.backend {
|
||||||
|
let err = Error::with_msg(format!(
|
||||||
|
"backend mismatch we {} requested {}",
|
||||||
|
node_config.node.backend, query.channel.backend
|
||||||
|
));
|
||||||
|
return Err(err);
|
||||||
|
}
|
||||||
let range = &query.range;
|
let range = &query.range;
|
||||||
let channel_config = read_local_config(&query.channel, &node_config.node).await?;
|
let channel_config = read_local_config(&query.channel, &node_config.node).await?;
|
||||||
let entry = extract_matching_config_entry(range, &channel_config);
|
let entry = extract_matching_config_entry(range, &channel_config);
|
||||||
@@ -249,6 +256,13 @@ pub fn pre_binned_bytes_for_http(
|
|||||||
query: &PreBinnedQuery,
|
query: &PreBinnedQuery,
|
||||||
) -> Result<PreBinnedValueByteStream, Error> {
|
) -> Result<PreBinnedValueByteStream, Error> {
|
||||||
info!("pre_binned_bytes_for_http {:?} {:?}", query, node_config.node);
|
info!("pre_binned_bytes_for_http {:?} {:?}", query, node_config.node);
|
||||||
|
if query.channel.backend != node_config.node.backend {
|
||||||
|
let err = Error::with_msg(format!(
|
||||||
|
"backend mismatch we {} requested {}",
|
||||||
|
node_config.node.backend, query.channel.backend
|
||||||
|
));
|
||||||
|
return Err(err);
|
||||||
|
}
|
||||||
let patch_node_ix = node_ix_for_patch(&query.patch, &query.channel, &node_config.node_config.cluster);
|
let patch_node_ix = node_ix_for_patch(&query.patch, &query.channel, &node_config.node_config.cluster);
|
||||||
if node_config.ix as u32 != patch_node_ix {
|
if node_config.ix as u32 != patch_node_ix {
|
||||||
Err(Error::with_msg(format!(
|
Err(Error::with_msg(format!(
|
||||||
|
|||||||
Vendored
+50
-31
@@ -8,12 +8,12 @@ use crate::streamlog::Streamlog;
|
|||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use err::Error;
|
use err::Error;
|
||||||
use futures_core::Stream;
|
use futures_core::Stream;
|
||||||
use futures_util::pin_mut;
|
|
||||||
use futures_util::{FutureExt, StreamExt};
|
use futures_util::{FutureExt, StreamExt};
|
||||||
use netpod::log::*;
|
use netpod::log::*;
|
||||||
use netpod::streamext::SCC;
|
use netpod::streamext::SCC;
|
||||||
use netpod::{BinnedRange, NodeConfigCached, PreBinnedPatchIterator, PreBinnedPatchRange};
|
use netpod::{BinnedRange, NodeConfigCached, PreBinnedPatchIterator, PreBinnedPatchRange};
|
||||||
use std::future::{ready, Future};
|
use std::future::{ready, Future};
|
||||||
|
use std::path::PathBuf;
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
use std::task::{Context, Poll};
|
use std::task::{Context, Poll};
|
||||||
|
|
||||||
@@ -53,6 +53,8 @@ pub struct PreBinnedValueStream {
|
|||||||
node_config: NodeConfigCached,
|
node_config: NodeConfigCached,
|
||||||
open_check_local_file: Option<Pin<Box<dyn Future<Output = Result<tokio::fs::File, std::io::Error>> + Send>>>,
|
open_check_local_file: Option<Pin<Box<dyn Future<Output = Result<tokio::fs::File, std::io::Error>> + Send>>>,
|
||||||
fut2: Option<Pin<Box<dyn Stream<Item = Result<PreBinnedItem, Error>> + Send>>>,
|
fut2: Option<Pin<Box<dyn Stream<Item = Result<PreBinnedItem, Error>> + Send>>>,
|
||||||
|
read_from_cache: bool,
|
||||||
|
cache_written: bool,
|
||||||
data_complete: bool,
|
data_complete: bool,
|
||||||
range_complete_observed: bool,
|
range_complete_observed: bool,
|
||||||
range_complete_emitted: bool,
|
range_complete_emitted: bool,
|
||||||
@@ -71,6 +73,8 @@ impl PreBinnedValueStream {
|
|||||||
node_config: node_config.clone(),
|
node_config: node_config.clone(),
|
||||||
open_check_local_file: None,
|
open_check_local_file: None,
|
||||||
fut2: None,
|
fut2: None,
|
||||||
|
read_from_cache: false,
|
||||||
|
cache_written: false,
|
||||||
data_complete: false,
|
data_complete: false,
|
||||||
range_complete_observed: false,
|
range_complete_observed: false,
|
||||||
range_complete_emitted: false,
|
range_complete_emitted: false,
|
||||||
@@ -196,9 +200,9 @@ impl Stream for PreBinnedValueStream {
|
|||||||
} else if let Some(item) = self.streamlog.pop() {
|
} else if let Some(item) = self.streamlog.pop() {
|
||||||
Ready(Some(Ok(PreBinnedItem::Log(item))))
|
Ready(Some(Ok(PreBinnedItem::Log(item))))
|
||||||
} else if let Some(fut) = &mut self.write_fut {
|
} else if let Some(fut) = &mut self.write_fut {
|
||||||
pin_mut!(fut);
|
match fut.poll_unpin(cx) {
|
||||||
match fut.poll(cx) {
|
|
||||||
Ready(item) => {
|
Ready(item) => {
|
||||||
|
self.cache_written = true;
|
||||||
self.write_fut = None;
|
self.write_fut = None;
|
||||||
match item {
|
match item {
|
||||||
Ok(()) => {
|
Ok(()) => {
|
||||||
@@ -220,6 +224,7 @@ impl Stream for PreBinnedValueStream {
|
|||||||
match item {
|
match item {
|
||||||
Ok(item) => {
|
Ok(item) => {
|
||||||
self.data_complete = true;
|
self.data_complete = true;
|
||||||
|
self.range_complete_observed = true;
|
||||||
Ready(Some(Ok(item)))
|
Ready(Some(Ok(item)))
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
@@ -230,33 +235,46 @@ impl Stream for PreBinnedValueStream {
|
|||||||
}
|
}
|
||||||
Pending => Pending,
|
Pending => Pending,
|
||||||
}
|
}
|
||||||
|
} else if self.range_complete_emitted {
|
||||||
|
self.completed = true;
|
||||||
|
Ready(None)
|
||||||
} else if self.data_complete {
|
} else if self.data_complete {
|
||||||
if self.range_complete_observed {
|
if self.cache_written {
|
||||||
if self.range_complete_emitted {
|
if self.range_complete_observed {
|
||||||
self.completed = true;
|
|
||||||
Ready(None)
|
|
||||||
} else {
|
|
||||||
let msg = format!(
|
|
||||||
"Write cache file\n{:?}\nN: {}\n\n\n",
|
|
||||||
self.query.patch,
|
|
||||||
self.values.ts1s.len()
|
|
||||||
);
|
|
||||||
self.streamlog.append(Level::INFO, msg);
|
|
||||||
let values = std::mem::replace(&mut self.values, MinMaxAvgScalarBinBatch::empty());
|
|
||||||
let fut = super::write_pb_cache_min_max_avg_scalar(
|
|
||||||
values,
|
|
||||||
self.query.patch.clone(),
|
|
||||||
self.query.agg_kind.clone(),
|
|
||||||
self.query.channel.clone(),
|
|
||||||
self.node_config.clone(),
|
|
||||||
);
|
|
||||||
self.write_fut = Some(Box::pin(fut));
|
|
||||||
self.range_complete_emitted = true;
|
self.range_complete_emitted = true;
|
||||||
Ready(Some(Ok(PreBinnedItem::RangeComplete)))
|
Ready(Some(Ok(PreBinnedItem::RangeComplete)))
|
||||||
|
} else {
|
||||||
|
self.completed = true;
|
||||||
|
Ready(None)
|
||||||
}
|
}
|
||||||
|
} else if self.read_from_cache {
|
||||||
|
self.cache_written = true;
|
||||||
|
continue 'outer;
|
||||||
} else {
|
} else {
|
||||||
self.completed = true;
|
match self.query.cache_usage {
|
||||||
Ready(None)
|
super::CacheUsage::Use | super::CacheUsage::Recreate => {
|
||||||
|
let msg = format!(
|
||||||
|
"Write cache file\n{:?}\nN: {}\n\n\n",
|
||||||
|
self.query.patch,
|
||||||
|
self.values.ts1s.len()
|
||||||
|
);
|
||||||
|
self.streamlog.append(Level::INFO, msg);
|
||||||
|
let values = std::mem::replace(&mut self.values, MinMaxAvgScalarBinBatch::empty());
|
||||||
|
let fut = super::write_pb_cache_min_max_avg_scalar(
|
||||||
|
values,
|
||||||
|
self.query.patch.clone(),
|
||||||
|
self.query.agg_kind.clone(),
|
||||||
|
self.query.channel.clone(),
|
||||||
|
self.node_config.clone(),
|
||||||
|
);
|
||||||
|
self.write_fut = Some(Box::pin(fut));
|
||||||
|
continue 'outer;
|
||||||
|
}
|
||||||
|
_ => {
|
||||||
|
self.cache_written = true;
|
||||||
|
continue 'outer;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
} else if let Some(fut) = self.fut2.as_mut() {
|
} else if let Some(fut) = self.fut2.as_mut() {
|
||||||
match fut.poll_next_unpin(cx) {
|
match fut.poll_next_unpin(cx) {
|
||||||
@@ -295,6 +313,7 @@ impl Stream for PreBinnedValueStream {
|
|||||||
self.open_check_local_file = None;
|
self.open_check_local_file = None;
|
||||||
match item {
|
match item {
|
||||||
Ok(file) => {
|
Ok(file) => {
|
||||||
|
self.read_from_cache = true;
|
||||||
let fut = super::read_pbv(file);
|
let fut = super::read_pbv(file);
|
||||||
self.read_cache_fut = Some(Box::pin(fut));
|
self.read_cache_fut = Some(Box::pin(fut));
|
||||||
continue 'outer;
|
continue 'outer;
|
||||||
@@ -323,17 +342,17 @@ impl Stream for PreBinnedValueStream {
|
|||||||
Pending => Pending,
|
Pending => Pending,
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
#[allow(unused_imports)]
|
|
||||||
use std::os::unix::fs::OpenOptionsExt;
|
|
||||||
let mut opts = std::fs::OpenOptions::new();
|
|
||||||
opts.read(true);
|
|
||||||
let cfd = CacheFileDesc {
|
let cfd = CacheFileDesc {
|
||||||
channel: self.query.channel.clone(),
|
channel: self.query.channel.clone(),
|
||||||
patch: self.query.patch.clone(),
|
patch: self.query.patch.clone(),
|
||||||
agg_kind: self.query.agg_kind.clone(),
|
agg_kind: self.query.agg_kind.clone(),
|
||||||
};
|
};
|
||||||
let path = cfd.path(&self.node_config);
|
use super::CacheUsage;
|
||||||
let fut = async { tokio::fs::OpenOptions::from(opts).open(path).await };
|
let path = match self.query.cache_usage {
|
||||||
|
CacheUsage::Use => cfd.path(&self.node_config),
|
||||||
|
_ => PathBuf::from("DOESNOTEXIST"),
|
||||||
|
};
|
||||||
|
let fut = async { tokio::fs::OpenOptions::new().read(true).open(path).await };
|
||||||
self.open_check_local_file = Some(Box::pin(fut));
|
self.open_check_local_file = Some(Box::pin(fut));
|
||||||
continue 'outer;
|
continue 'outer;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -62,6 +62,7 @@ pub async fn gen_test_data() -> Result<(), Error> {
|
|||||||
split: i1,
|
split: i1,
|
||||||
data_base_path: data_base_path.join(format!("node{:02}", i1)),
|
data_base_path: data_base_path.join(format!("node{:02}", i1)),
|
||||||
ksprefix: ksprefix.clone(),
|
ksprefix: ksprefix.clone(),
|
||||||
|
backend: "testbackend".into(),
|
||||||
};
|
};
|
||||||
ensemble.nodes.push(node);
|
ensemble.nodes.push(node);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -105,6 +105,7 @@ pub struct Node {
|
|||||||
pub split: u32,
|
pub split: u32,
|
||||||
pub data_base_path: PathBuf,
|
pub data_base_path: PathBuf,
|
||||||
pub ksprefix: String,
|
pub ksprefix: String,
|
||||||
|
pub backend: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||||
|
|||||||
@@ -60,6 +60,7 @@ fn simple_fetch() {
|
|||||||
data_base_path: err::todoval(),
|
data_base_path: err::todoval(),
|
||||||
ksprefix: "daq_swissfel".into(),
|
ksprefix: "daq_swissfel".into(),
|
||||||
split: 0,
|
split: 0,
|
||||||
|
backend: "testbackend".into(),
|
||||||
};
|
};
|
||||||
let query = netpod::AggQuerySingleChannel {
|
let query = netpod::AggQuerySingleChannel {
|
||||||
channel_config: ChannelConfig {
|
channel_config: ChannelConfig {
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ fn test_cluster() -> Cluster {
|
|||||||
data_base_path: format!("../tmpdata/node{:02}", id).into(),
|
data_base_path: format!("../tmpdata/node{:02}", id).into(),
|
||||||
ksprefix: "ks".into(),
|
ksprefix: "ks".into(),
|
||||||
split: id,
|
split: id,
|
||||||
|
backend: "testbackend".into(),
|
||||||
})
|
})
|
||||||
.collect();
|
.collect();
|
||||||
Cluster {
|
Cluster {
|
||||||
|
|||||||
Reference in New Issue
Block a user