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
94 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,95 @@ | ||
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>; | ||
/// Represents a stream reader that deserializes JSON objects from a reader. | ||
pub struct StreamReader<T> { | ||
reader: BufReader<Box<dyn Read>>, | ||
_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: DeserializeOwned> StreamReader<T> { | ||
/// Creates a new `StreamReader` from a boxed reader. | ||
pub fn new(reader: Box<dyn Read>) -> 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> Iterator for StreamReader<T> | ||
where | ||
T: DeserializeOwned, | ||
{ | ||
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>; | ||
} | ||
|
||
impl RequestExt for Response { | ||
fn into_stream<T: DeserializeOwned>(self) -> StreamReader<T> { | ||
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