use futures_util::StreamExt;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::RwLock;
use tokio_tungstenite::connect_async;
use crate::Config;
use crate::actor::{Actor, ActorContext, Addr};
use crate::adapters::ws_conn::WsConn;
use crate::message::Message;
use async_trait::async_trait;
use log::{debug, info};
use tokio::time::{Duration, sleep};
#[derive(Clone)]
pub struct OutgoingWebsocketManager {
config: Config,
clients: Arc<RwLock<HashMap<String, Addr>>>,
urls: Vec<String>,
}
impl OutgoingWebsocketManager {
pub fn new(config: Config, urls: Vec<String>) -> Self {
OutgoingWebsocketManager {
urls,
clients: Arc::new(RwLock::new(HashMap::new())),
config,
}
}
pub async fn connected_count(&self) -> usize {
self.clients.read().await.len()
}
pub fn urls(&self) -> &[String] {
&self.urls
}
}
#[async_trait]
impl Actor for OutgoingWebsocketManager {
async fn pre_start(&mut self, ctx: &ActorContext) {
info!("OutgoingWebsocketManager starting");
for url in self.urls.iter() {
loop {
if self.clients.read().await.contains_key(url) {
debug!("already connected to {}", url);
break;
}
debug!("attempting WebSocket connect to {}", url);
let result = connect_async(url).await;
if let Ok((socket, _)) = result {
let (sender, receiver) = socket.split();
let client = WsConn::new(sender, receiver, self.config.allow_public_space);
let addr = ctx.start_actor(Box::new(client));
self.clients.write().await.insert(url.clone(), addr);
debug!("connected to {}", url);
break;
}
debug!("connect to {} failed, retrying in 200ms", url);
sleep(Duration::from_millis(200)).await;
}
}
}
fn subscribe_to_everything(&self) -> bool {
true
}
async fn handle(&mut self, message: Message, _ctx: &ActorContext) {
let snapshot: Vec<Addr> = self.clients.read().await.values().cloned().collect();
for client in snapshot {
let _ = client.send(message.clone());
}
}
async fn stopping(&mut self, _ctx: &ActorContext) {
let count = self.clients.read().await.len();
info!(
"OutgoingWebsocketManager stopping — {} outgoing connections",
count
);
self.clients.write().await.clear();
}
}