forked from oreoslabs/oreowallet-mono
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
add rpc streaming from ironfish node
- Loading branch information
Showing
5 changed files
with
93 additions
and
28 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,33 +1,94 @@ | ||
use std::io::Read; | ||
use std::io::{BufRead, BufReader, Read}; | ||
use std::marker::PhantomData; | ||
|
||
use oreo_errors::OreoError; | ||
use serde::de::DeserializeOwned; | ||
use serde::Deserialize; | ||
use ureq::Response; | ||
|
||
use crate::rpc_abi::RpcResponseStream; | ||
|
||
pub trait RequestExt { | ||
fn collect_stream<T: DeserializeOwned>(self) -> Result<Vec<T>, OreoError>; | ||
pub struct StreamReader<T, R> { | ||
reader: BufReader<R>, | ||
_marker: PhantomData<T>, | ||
} | ||
|
||
impl RequestExt for Response { | ||
fn collect_stream<T: DeserializeOwned>(self) -> Result<Vec<T>, OreoError> { | ||
let reader = self.into_reader(); | ||
let mut buffered = std::io::BufReader::new(reader); | ||
let mut items = Vec::new(); | ||
let mut response_str = String::new(); | ||
buffered.read_to_string(&mut response_str).map_err(|e| OreoError::InternalRpcError(e.to_string()))?; | ||
let lines = response_str.split('\x0c').collect::<Vec<&str>>(); | ||
|
||
// Get rid of status code | ||
for line in lines[0..lines.len()-1].into_iter() { | ||
let line = *line; // Dereference to get &str | ||
if !line.trim().is_empty() { | ||
let item: RpcResponseStream<T> = serde_json::from_str(line) | ||
.map_err(|e| OreoError::InternalRpcError(e.to_string()))?; | ||
items.push(item.data); | ||
#[derive(Deserialize)] | ||
#[serde(untagged)] | ||
enum ResponseItem<T> { | ||
Data(RpcResponseStream<T>), | ||
Status { status: u16 }, | ||
} | ||
|
||
impl<T, R> StreamReader<T, R> where T: DeserializeOwned, R: Read { | ||
pub fn new(reader: R) -> Self { | ||
Self { | ||
reader: BufReader::new(reader), | ||
_marker: PhantomData, | ||
} | ||
} | ||
|
||
fn read_item(&mut self) -> Result<Option<Vec<u8>>, OreoError> { | ||
let mut item = Vec::new(); | ||
loop { | ||
let bytes_read = self.reader.read_until(b'\x0c', &mut item) | ||
.map_err(|e| OreoError::RpcStreamError(e.to_string()))?; | ||
if bytes_read == 0 { | ||
break; | ||
} | ||
if item.last() == Some(&b'\x0c') { | ||
item.pop(); | ||
break; | ||
} | ||
} | ||
match item.len() { | ||
0 => Ok(None), | ||
_ => Ok(Some(item)) | ||
} | ||
} | ||
|
||
|
||
/// Parses a chunk of data into a `ResponseItem<T>`. | ||
fn parse_item(&self, chunk: &[u8]) -> Result<Option<T>, OreoError> { | ||
match serde_json::from_slice::<ResponseItem<T>>(chunk) { | ||
Ok(ResponseItem::Data(item)) => Ok(Some(item.data)), | ||
Ok(ResponseItem::Status { status: 200 }) => Ok(None), | ||
Ok(ResponseItem::Status { status }) => Err(OreoError::RpcStreamError(format!("Received error status: {}", status))), | ||
Err(e) => { | ||
let err_str = format!("Failed to parse JSON object: {:?}", e); | ||
Err(OreoError::RpcStreamError(err_str)) | ||
} | ||
} | ||
Ok(items) | ||
} | ||
} | ||
} | ||
|
||
impl<T, R> Iterator for StreamReader<T, R> | ||
where | ||
T: DeserializeOwned, | ||
R: Read, | ||
{ | ||
type Item = Result<T, OreoError>; | ||
|
||
fn next(&mut self) -> Option<Self::Item> { | ||
match self.read_item() { | ||
Ok(Some(chunk)) => match self.parse_item(&chunk) { | ||
Ok(Some(data)) => Some(Ok(data)), | ||
Ok(None) => None, // End of stream | ||
Err(e) => Some(Err(e)), | ||
}, | ||
Ok(None) => None, // EOF reached | ||
Err(e) => Some(Err(e)), | ||
} | ||
} | ||
} | ||
|
||
pub trait RequestExt { | ||
fn into_stream<T: DeserializeOwned>(self) -> StreamReader<T, Box<dyn Read + Send + Sync + 'static>>; | ||
} | ||
|
||
impl RequestExt for Response { | ||
fn into_stream<T: DeserializeOwned>(self) -> StreamReader<T, Box<dyn Read + Send + Sync + 'static>> { | ||
let reader = self.into_reader(); | ||
StreamReader::new(Box::new(reader)) | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters