use crate::error::Error;
use futures::{
future::BoxFuture,
stream::{BoxStream, StreamExt, TryStreamExt},
};
#[cfg(feature = "12-47-0")]
use misskey_api::model::channel::Channel;
use misskey_api::model::{antenna::Antenna, note::Note, query::Query, user_list::UserList};
use misskey_api::{
streaming::{self, channel},
EntityRef,
};
use misskey_core::streaming::StreamingClient;
#[allow(clippy::type_complexity)]
pub trait StreamingClientExt: StreamingClient + Sync {
fn subscribe_note(
&self,
note: impl EntityRef<Note>,
) -> BoxFuture<
Result<
BoxStream<Result<streaming::note::NoteUpdateEvent, Error<Self::Error>>>,
Error<Self::Error>,
>,
> {
let note_id = note.entity_ref().to_string();
Box::pin(async move {
Ok(self
.subnote(note_id)
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.boxed())
})
}
fn main_stream(
&self,
) -> BoxFuture<
Result<
BoxStream<Result<channel::main::MainStreamEvent, Error<Self::Error>>>,
Error<Self::Error>,
>,
> {
Box::pin(async move {
Ok(self
.channel(channel::main::Request::default())
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.boxed())
})
}
fn home_timeline(
&self,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::home_timeline::{HomeTimelineEvent, Request};
Box::pin(async move {
Ok(self
.channel(Request::default())
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|HomeTimelineEvent::Note(note)| note)
.boxed())
})
}
fn local_timeline(
&self,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::local_timeline::{LocalTimelineEvent, Request};
Box::pin(async move {
Ok(self
.channel(Request::default())
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|LocalTimelineEvent::Note(note)| note)
.boxed())
})
}
fn social_timeline(
&self,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::hybrid_timeline::{HybridTimelineEvent, Request};
Box::pin(async move {
Ok(self
.channel(Request::default())
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|HybridTimelineEvent::Note(note)| note)
.boxed())
})
}
fn global_timeline(
&self,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::global_timeline::{GlobalTimelineEvent, Request};
Box::pin(async move {
Ok(self
.channel(Request::default())
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|GlobalTimelineEvent::Note(note)| note)
.boxed())
})
}
fn hashtag_timeline(
&self,
query: impl Into<Query<String>>,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::hashtag::{HashtagEvent, Request};
let q = query.into();
Box::pin(async move {
Ok(self
.channel(Request { q })
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|HashtagEvent::Note(note)| note)
.boxed())
})
}
fn antenna_timeline(
&self,
antenna: impl EntityRef<Antenna>,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::antenna::{AntennaStreamEvent, Request};
let antenna_id = antenna.entity_ref();
Box::pin(async move {
Ok(self
.channel(Request { antenna_id })
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|AntennaStreamEvent::Note(note)| note)
.boxed())
})
}
#[cfg(feature = "12-47-0")]
#[cfg_attr(docsrs, doc(cfg(feature = "12-47-0")))]
fn channel_timeline(
&self,
channel: impl EntityRef<Channel>,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::channel::{ChannelEvent, Request};
let channel_id = channel.entity_ref();
Box::pin(async move {
Ok(self
.channel(Request { channel_id })
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.map_ok(|ChannelEvent::Note(note)| note)
.boxed())
})
}
fn user_list_timeline(
&self,
list: impl EntityRef<UserList>,
) -> BoxFuture<Result<BoxStream<Result<Note, Error<Self::Error>>>, Error<Self::Error>>> {
use channel::user_list::{Request, UserListEvent};
let list_id = list.entity_ref();
Box::pin(async move {
Ok(self
.channel(Request { list_id })
.await
.map_err(Error::Client)?
.map_err(Error::Client)
.try_filter_map(|event| async move {
if let UserListEvent::Note(note) = event {
Ok(Some(note))
} else {
Ok(None)
}
})
.boxed())
})
}
}
impl<C: StreamingClient + Sync> StreamingClientExt for C {}