pub mod types;
use futures::channel::mpsc::{unbounded, UnboundedReceiver, UnboundedSender};
use futures::{SinkExt, StreamExt};
use serde::Serialize;
use serde_json::{json, Value};
use tokio::spawn;
use self::types::{Event, EventCapsule};
use crate::errors::Result;
use crate::manager::types::ShardManagerMessage;
use crate::manager::{ManagerOptions, ShardManager};
pub struct LeapOptions<'a> {
pub token: Option<&'a str>,
pub project: &'a str,
pub ws_url: &'a str,
}
impl Default for LeapOptions<'_> {
fn default() -> Self {
Self {
token: None,
project: "",
ws_url: "wss://leap.hop.io/ws",
}
}
}
pub struct LeapEdge {
manager_tx: UnboundedSender<ShardManagerMessage>,
leap_rx: UnboundedReceiver<Event>,
}
impl LeapEdge {
pub async fn new(options: LeapOptions<'_>) -> Result<Self> {
let (leap_tx, leap_rx) = unbounded();
let mut manager = ShardManager::new(ManagerOptions {
project: options.project,
ws_url: options.ws_url,
token: options.token,
event_tx: leap_tx.clone(),
})
.await?;
let manager_tx = manager.get_manager_tx();
spawn(async move {
if let Err(why) = manager.run().await {
log::debug!("[Manager] Stopped: {why:?}");
} else {
log::debug!("[Manager] Stopped");
}
});
Ok(Self {
manager_tx,
leap_rx,
})
}
#[inline]
pub async fn send_service_message<D>(&mut self, message: D) -> Result<()>
where
D: Serialize,
{
self.manager_tx
.send(ShardManagerMessage::Json(json!({
"op": 0,
"d": message,
})))
.await?;
Ok(())
}
#[inline]
pub async fn channel_subscribe(&mut self, channel: &str) -> Result<()> {
self.channel_subscribe_with_data(channel, Value::Null).await
}
#[inline]
pub async fn channel_subscribe_with_data<D>(&mut self, channel: &str, data: D) -> Result<()>
where
D: Serialize,
{
self.send_service_message(&Event::Subscribe(EventCapsule {
channel: Some(channel.to_string()),
data: serde_json::to_value(data)?,
unicast: false,
}))
.await
}
#[inline]
pub async fn listen(&mut self) -> Option<Event> {
self.leap_rx.next().await
}
pub async fn close(&mut self) {
self.manager_tx.send(ShardManagerMessage::Close).await.ok();
}
}