use std::sync::Arc;
use futures::{Stream, StreamExt};
use futures::stream::{iter, unfold};
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex;
use crate::action::{CastAction, LinkAction, LinkInfo, ReactionAction};
use crate::error::FatlineError;
use crate::HubService;
use crate::posts::{Cast, CastId, Embed, Mention, Parent, PostService};
use crate::proto::{FidRequest, FidTimestampRequest, FidsRequest, Message, OnChainEvent, OnChainEventRequest, LinkBody, MessageData, MessageType, ReactionType, CastsByParentRequest, ReactionsByTargetRequest};
use crate::proto::message_data::Body;
use crate::proto::reaction_body::Target::{TargetCastId, TargetUrl};
use crate::utils;
pub(crate) type AccumulationFidInput<'a> = (u64, bool, Option<u32>, Arc<Mutex<&'a mut HubService>>, Option<Vec<u8>>, usize);
pub(crate) type AccumulationFidResponse<'a, T> = Option<(Result<Vec<T>,FatlineError>, AccumulationFidInput<'a>)>;
pub(crate) type AccumulationParentInput<'a> = (Parent, bool, Option<u32>, Arc<Mutex<&'a mut HubService>>, Option<Vec<u8>>, usize);
pub(crate) type AccumulationParentResponse<'a, T> = Option<(Result<Vec<T>, FatlineError>, AccumulationParentInput<'a>)>;
pub(crate) type AccumulationGlobalInput<'a> = (bool, Option<u32>, Arc<Mutex<&'a mut HubService>>, Option<Vec<u8>>, usize);
pub(crate) type AccumulationGlobalResponse<'a, T> = Option<(Result<Vec<T>, FatlineError>, AccumulationGlobalInput<'a>)>;
pub trait StreamService {
fn get_all_cast_messages_by_fid_stream(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=CastAction>;
fn get_all_cast_messages_by_parent(&mut self, parent: Parent, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=CastAction>;
fn get_all_cast_reacts_by_parent(&mut self, parent: Parent, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=ReactionAction>;
fn get_all_onchain_signers(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=Result<Vec<OnChainEvent>, FatlineError>>;
fn get_all_onchain_fid_events(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item = Result<Vec<OnChainEvent>, FatlineError>>;
fn get_all_fids(&mut self, reverse: bool, page_size: Option<u32>) -> impl Stream<Item = Result<Vec<u64>, FatlineError>>;
fn get_all_profile_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=Result<Vec<Message>, FatlineError>>;
fn get_all_link_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=LinkAction>;
fn get_all_react_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=ReactionAction>;
}
impl StreamService for HubService {
fn get_all_cast_messages_by_fid_stream(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=CastAction> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_cast_messages)
.filter_map(|res| async { res.ok() })
.flat_map(|messages| {
let mapped = messages.iter().cloned().filter_map(utils::cast_action_from_message).collect::<Vec<_>>();
iter(mapped)
})
}
fn get_all_cast_messages_by_parent(&mut self, parent: Parent, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=CastAction> {
unfold((parent, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_cast_messages_by_parent)
.filter_map(|res| async { res.ok() })
.flat_map(|messages| {
let mapped = messages.iter().cloned().filter_map(utils::cast_action_from_message).collect::<Vec<_>>();
iter(mapped)
})
}
fn get_all_cast_reacts_by_parent(&mut self, parent: Parent, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=ReactionAction> {
unfold((parent, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_cast_reacts_by_parent)
.filter_map(|res| async { res.ok() })
.flat_map(|messages| {
let mapped = messages.iter().cloned().filter_map(utils::reaction_from_message).collect::<Vec<_>>();
iter(mapped)
})
}
fn get_all_onchain_signers(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=Result<Vec<OnChainEvent>, FatlineError>> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_signer_events).boxed()
}
fn get_all_onchain_fid_events(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=Result<Vec<OnChainEvent>, FatlineError>> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_fid_events)
}
fn get_all_fids(&mut self, reverse: bool, page_size: Option<u32>) -> impl Stream<Item = Result<Vec<u64>, FatlineError>> {
unfold((reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_fids).boxed()
}
fn get_all_profile_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=Result<Vec<Message>, FatlineError>> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_profile_events).boxed()
}
fn get_all_link_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=LinkAction> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_link_events).filter_map(|res| async {res.ok()})
.flat_map(|messages| {
let mapped = messages.iter().cloned()
.filter_map(utils::link_from_message)
.flat_map(|iter| iter)
.collect::<Vec<_>>();
iter(mapped)
})
}
fn get_all_react_messages(&mut self, fid: u64, reverse: bool, page_size: Option<u32>) -> impl Stream<Item=ReactionAction> {
unfold((fid, reverse, page_size, Arc::new(Mutex::new(self)), None, 0), accumulate_all_react_events).filter_map(|res| async {res.ok()})
.flat_map(|messages| {
let mapped = messages.iter().cloned()
.filter_map(utils::reaction_from_message)
.collect::<Vec<_>>();
iter(mapped)
})
}
}
fn should_return_early(previous_size: usize, page: &Option<Vec<u8>>) -> bool {
previous_size > 0 && (page.is_none() || page == &Some(vec![]))
}
pub fn flat_map_result<T>(result: Result<Vec<T>, FatlineError>) -> impl Stream<Item = Result<T,FatlineError>>
where T: Clone {
let items = match result {
Ok(items) => items.iter().cloned().map(Ok).collect(),
Err(e) => vec![Err(e)]
};
iter(items)
}
async fn accumulate_all_cast_reacts_by_parent<'a>((parent, reverse, page_size, client_arc, page, previous_size): AccumulationParentInput<'a>) -> AccumulationParentResponse<'a, Message> {
let mut client= client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = ReactionsByTargetRequest {
reaction_type: None,
page_size,
page_token: page.clone(),
reverse: Some(reverse),
target: Some(parent.clone().into()),
};
match (*client).get_reactions_by_target(request).await {
Ok(response) => {
let inner = response.into_inner();
let messages = inner.messages;
let size = messages.len();
if messages.is_empty() {
None
} else {
Some((Ok(messages), (parent, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (parent, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_cast_messages_by_parent<'a>((parent, reverse, page_size, client_arc, page, previous_size): AccumulationParentInput<'a>) -> AccumulationParentResponse<'a, Message> {
let mut client= client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = CastsByParentRequest {
page_size,
page_token: page.clone(),
reverse: Some(reverse),
parent: Some(parent.clone().into()),
};
match (*client).get_casts_by_parent(request).await {
Ok(response) => {
let inner = response.into_inner();
let messages = inner.messages;
let size = messages.len();
if messages.is_empty() {
None
} else {
Some((Ok(messages), (parent, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (parent, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_cast_messages<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a,Message> {
let mut client= client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidTimestampRequest {
fid,
page_size,
page_token: page.clone(),
reverse: Some(reverse),
start_timestamp: None,
stop_timestamp: None,
};
match (*client).get_all_cast_messages_by_fid(request).await {
Ok(response) => {
let inner = response.into_inner();
let messages = inner.messages;
let size = messages.len();
if messages.is_empty() {
None
} else {
Some((Ok(messages), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_link_events<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a, Message> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidTimestampRequest {
fid,
page_size,
reverse: Some(reverse),
start_timestamp: None,
page_token: page.clone(),
stop_timestamp: None,
};
match (*client).get_all_link_messages_by_fid(request).await {
Ok(response) => {
let inner = response.into_inner();
let messages = inner.messages;
let size = messages.len();
if messages.is_empty() {
None
} else {
Some((Ok(messages), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_react_events<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a, Message> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidTimestampRequest {
fid,
page_size,
reverse: Some(reverse),
start_timestamp: None,
page_token: page.clone(),
stop_timestamp: None,
};
match (*client).get_all_reaction_messages_by_fid(request).await {
Ok(response) => {
let inner = response.into_inner();
let messages = inner.messages;
let size = messages.len();
if messages.is_empty() {
None
} else {
Some((Ok(messages), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_profile_events<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a, Message> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidTimestampRequest {
fid,
page_size,
reverse: Some(reverse),
start_timestamp: None,
page_token: page.clone(),
stop_timestamp: None,
};
match (*client).get_all_user_data_messages_by_fid(request).await {
Ok(response) => {
let inner = response.into_inner();
let events = inner.messages;
let size = events.len();
if events.is_empty() {
None
} else {
Some((Ok(events), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_signer_events<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a, OnChainEvent> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidRequest {
fid,
page_size,
reverse: Some(reverse),
page_token: page.clone()
};
match (*client).get_on_chain_signers_by_fid(request).await {
Ok(response) => {
let inner = response.into_inner();
let events = inner.events;
let size = events.len();
if events.is_empty() {
None
} else {
Some((Ok(events), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_fid_events<'a>((fid, reverse, page_size, client_arc, page, previous_size): AccumulationFidInput<'a>) -> AccumulationFidResponse<'a,OnChainEvent> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = OnChainEventRequest {
fid,
event_type: 0,
page_size,
page_token: page.clone(),
reverse: Some(reverse)
};
match (*client).get_on_chain_events(request).await {
Ok(response) => {
let inner = response.into_inner();
let events: Vec<_> = inner.events;
let size = events.len();
if events.is_empty() {
None
} else {
Some((Ok(events), (fid, reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => {
Some((Err(FatlineError::ClientError(e)), (fid, reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
}
async fn accumulate_all_fids<'a>((reverse, page_size, client_arc, page, previous_size): AccumulationGlobalInput<'a>) -> AccumulationGlobalResponse<'a,u64> {
let mut client = client_arc.lock().await;
if should_return_early(previous_size, &page) {
return None;
}
let request = FidsRequest {
reverse: Some(reverse),
page_token: page.clone(),
page_size,
};
match (*client).get_fids(request).await {
Ok(response) => {
let inner = response.into_inner();
let fids = inner.fids;
let size = fids.len();
if fids.is_empty() {
None
} else {
Some((Ok(fids), (reverse, page_size, client_arc.clone(), inner.next_page_token, size)))
}
},
Err(e) => Some((Err(FatlineError::ClientError(e)), (reverse, page_size, client_arc.clone(), page, previous_size)))
}
}
#[cfg(test)]
mod test {
use std::pin::pin;
use futures::StreamExt;
use tonic::transport::Uri;
use crate::error::FatlineError;
use crate::HubService;
use crate::proto::Message;
use crate::stream::{flat_map_result, StreamService};
async fn print_events_result(result: Result<Message, FatlineError>) {
match result {
Ok(event) => println!("received {event:?}"),
Err(err) => eprintln!("{err:?}")
};
}
#[tokio::test]
async fn test_fid_retrieval() {
let hub_url = dotenvy::var("HUB_URL").expect("No HUB_URL present");
let uri = hub_url.parse::<Uri>().expect("Hub URL doesn't parse to valid URI");
let mut client = HubService::connect(uri).await.expect("couldn't connect for events");
let mut event_client = client.clone();
let mut fid_stream = client.get_all_fids(false, Some(10)).flat_map(flat_map_result);
while let Some(Ok(fid)) = fid_stream.next().await {
let mut cast_stream = event_client.get_all_cast_messages_by_fid_stream(fid, true, Some(10));
let mut cast_stream = pin!(cast_stream);
while let Some(message) = cast_stream.next().await {
println!("next message {message:?}");
}
}
}
}