Factor usage of common error type more
This commit is contained in:
@@ -38,7 +38,7 @@ impl<TBT> FetchedPreBinned<TBT> {
|
||||
let mut url = Url::parse(&format!("http://{}:{}/api/4/prebinned", node.host, node.port))?;
|
||||
query.append_to_url(&mut url);
|
||||
let ret = Self {
|
||||
uri: Uri::from_str(&url.to_string())?,
|
||||
uri: Uri::from_str(&url.to_string()).map_err(Error::from_string)?,
|
||||
resfut: None,
|
||||
res: None,
|
||||
errored: false,
|
||||
@@ -115,7 +115,7 @@ where
|
||||
Err(e) => {
|
||||
error!("PreBinnedValueStream error in stream {:?}", e);
|
||||
self.errored = true;
|
||||
Ready(Some(Err(e.into())))
|
||||
Ready(Some(Err(Error::from_string(e))))
|
||||
}
|
||||
},
|
||||
Pending => Pending,
|
||||
@@ -133,7 +133,7 @@ where
|
||||
}
|
||||
Err(e) => {
|
||||
self.errored = true;
|
||||
Ready(Some(Err(e.into())))
|
||||
Ready(Some(Err(Error::from_string(e))))
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
@@ -209,7 +209,8 @@ where
|
||||
Ok::<_, Error>(enc.len())
|
||||
}
|
||||
})
|
||||
.await??;
|
||||
.await
|
||||
.map_err(Error::from_string)??;
|
||||
tokio::fs::rename(&tmp_path, &path).await?;
|
||||
let ts2 = Instant::now();
|
||||
let ret = WrittenPbCache {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use super::paths;
|
||||
use bytes::BytesMut;
|
||||
use err::Error;
|
||||
use err::{ErrStr, Error};
|
||||
use futures_util::StreamExt;
|
||||
use netpod::log::*;
|
||||
use netpod::{ChannelConfig, NanoRange, Nanos, Node};
|
||||
@@ -251,7 +251,7 @@ async fn open_files_inner(
|
||||
"----- open_files_inner giving OpenedFileSet with {} files",
|
||||
h.files.len()
|
||||
);
|
||||
chtx.send(Ok(h)).await?;
|
||||
chtx.send(Ok(h)).await.errstr()?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -361,7 +361,7 @@ async fn open_expanded_files_inner(
|
||||
"----- open_expanded_files_inner giving OpenedFileSet with {} files",
|
||||
h.files.len()
|
||||
);
|
||||
chtx.send(Ok(h)).await?;
|
||||
chtx.send(Ok(h)).await.errstr()?;
|
||||
if found_pre {
|
||||
p1 += 1;
|
||||
break;
|
||||
@@ -383,7 +383,7 @@ async fn open_expanded_files_inner(
|
||||
}
|
||||
}
|
||||
let h = OpenedFileSet { timebin: tb, files: a };
|
||||
chtx.send(Ok(h)).await?;
|
||||
chtx.send(Ok(h)).await.errstr()?;
|
||||
p1 += 1;
|
||||
}
|
||||
} else {
|
||||
@@ -417,7 +417,7 @@ fn expanded_file_list() {
|
||||
array: false,
|
||||
compression: false,
|
||||
};
|
||||
let cluster = taskrun::test_cluster();
|
||||
let cluster = netpod::test_cluster();
|
||||
let task = async move {
|
||||
let mut paths = vec![];
|
||||
let mut files = open_expanded_files(&range, &channel_config, cluster.nodes[0].clone());
|
||||
|
||||
@@ -229,7 +229,7 @@ mod test {
|
||||
array: false,
|
||||
compression: false,
|
||||
};
|
||||
let cluster = taskrun::test_cluster();
|
||||
let cluster = netpod::test_cluster();
|
||||
let node = cluster.nodes[nodeix].clone();
|
||||
let buffer_size = 512;
|
||||
let event_chunker_conf = EventChunkerConf {
|
||||
|
||||
Reference in New Issue
Block a user