use std::sync::atomic::{AtomicU64, Ordering};
use crate::client::AriClient;
use crate::error::Result;
use crate::event::{AriEvent, AriMessage, Bridge, Channel};
use crate::resources::bridge::BridgeHandle;
use crate::resources::channel::{ChannelHandle, OriginateParams};
use asterisk_rs_core::event::FilteredSubscription;
static PENDING_COUNTER: AtomicU64 = AtomicU64::new(1);
fn generate_pending_id(prefix: &str) -> String {
let id = PENDING_COUNTER.fetch_add(1, Ordering::Relaxed);
format!("{prefix}-pending-{id}")
}
fn event_matches_channel_id(event: &AriEvent, id: &str) -> bool {
match event {
AriEvent::StasisStart { channel, .. }
| AriEvent::StasisEnd { channel }
| AriEvent::ChannelCreated { channel }
| AriEvent::ChannelDestroyed { channel, .. }
| AriEvent::ChannelStateChange { channel }
| AriEvent::ChannelDtmfReceived { channel, .. }
| AriEvent::ChannelHangupRequest { channel }
| AriEvent::ChannelCallerId { channel, .. }
| AriEvent::ChannelConnectedLine { channel }
| AriEvent::ChannelDialplan { channel, .. }
| AriEvent::ChannelHold { channel, .. }
| AriEvent::ChannelUnhold { channel }
| AriEvent::ChannelTalkingStarted { channel }
| AriEvent::ChannelTalkingFinished { channel, .. }
| AriEvent::ChannelToneDetected { channel }
| AriEvent::ChannelTransfer { channel, .. }
| AriEvent::ChannelEnteredBridge { channel, .. }
| AriEvent::ChannelLeftBridge { channel, .. }
| AriEvent::ApplicationMoveFailed { channel, .. } => channel.id == id,
AriEvent::Dial {
peer,
caller,
forwarded,
..
} => {
peer.id == id
|| caller.as_ref().is_some_and(|c| c.id == id)
|| forwarded.as_ref().is_some_and(|c| c.id == id)
}
AriEvent::BridgeBlindTransfer {
channel,
transferee,
replace_channel,
..
} => {
channel.id == id
|| transferee.as_ref().is_some_and(|c| c.id == id)
|| replace_channel.as_ref().is_some_and(|c| c.id == id)
}
_ => false,
}
}
fn event_matches_bridge_id(event: &AriEvent, id: &str) -> bool {
match event {
AriEvent::BridgeCreated { bridge }
| AriEvent::BridgeDestroyed { bridge }
| AriEvent::ChannelEnteredBridge { bridge, .. }
| AriEvent::ChannelLeftBridge { bridge, .. }
| AriEvent::BridgeVideoSourceChanged { bridge, .. } => bridge.id == id,
AriEvent::BridgeMerged {
bridge,
bridge_from,
} => bridge.id == id || bridge_from.id == id,
AriEvent::BridgeBlindTransfer { bridge, .. } => bridge.as_ref().is_some_and(|b| b.id == id),
AriEvent::BridgeAttendedTransfer {
transferer_first_leg_bridge,
transferer_second_leg_bridge,
destination_threeway_bridge,
..
} => {
transferer_first_leg_bridge
.as_ref()
.is_some_and(|b| b.id == id)
|| transferer_second_leg_bridge
.as_ref()
.is_some_and(|b| b.id == id)
|| destination_threeway_bridge
.as_ref()
.is_some_and(|b| b.id == id)
}
_ => false,
}
}
fn event_playback_id(event: &AriEvent) -> Option<&str> {
match event {
AriEvent::PlaybackStarted { playback }
| AriEvent::PlaybackFinished { playback }
| AriEvent::PlaybackContinuing { playback } => Some(&playback.id),
_ => None,
}
}
#[derive(Debug)]
pub struct PendingChannel {
id: String,
client: AriClient,
events: FilteredSubscription<AriMessage>,
}
impl PendingChannel {
pub(crate) fn new(client: AriClient) -> Self {
let id = generate_pending_id("channel");
let filter_id = id.clone();
let events =
client.subscribe_filtered(move |msg| event_matches_channel_id(&msg.event, &filter_id));
Self { id, client, events }
}
pub fn id(&self) -> &str {
&self.id
}
pub async fn originate(
self,
mut params: OriginateParams,
) -> Result<(ChannelHandle, FilteredSubscription<AriMessage>)> {
params.channel_id = Some(self.id.clone());
let channel: Channel = self.client.post("/channels", ¶ms).await?;
let handle = ChannelHandle::new(channel.id, self.client);
Ok((handle, self.events))
}
pub fn events_mut(&mut self) -> &mut FilteredSubscription<AriMessage> {
&mut self.events
}
}
#[derive(Debug)]
pub struct PendingBridge {
id: String,
client: AriClient,
events: FilteredSubscription<AriMessage>,
}
impl PendingBridge {
pub(crate) fn new(client: AriClient) -> Self {
let id = generate_pending_id("bridge");
let filter_id = id.clone();
let events =
client.subscribe_filtered(move |msg| event_matches_bridge_id(&msg.event, &filter_id));
Self { id, client, events }
}
pub fn id(&self) -> &str {
&self.id
}
pub async fn create(
self,
bridge_type: &str,
) -> Result<(BridgeHandle, FilteredSubscription<AriMessage>)> {
let bridge: Bridge = self
.client
.post(
"/bridges",
&serde_json::json!({ "bridgeId": self.id, "type": bridge_type }),
)
.await?;
let handle = BridgeHandle::new(bridge.id, self.client);
Ok((handle, self.events))
}
pub fn events_mut(&mut self) -> &mut FilteredSubscription<AriMessage> {
&mut self.events
}
}
#[derive(Debug)]
pub struct PendingPlayback {
id: String,
events: FilteredSubscription<AriMessage>,
}
impl PendingPlayback {
pub(crate) fn new(client: &AriClient) -> Self {
let id = generate_pending_id("playback");
let filter_id = id.clone();
let events = client.subscribe_filtered(move |msg| {
event_playback_id(&msg.event).is_some_and(|pb_id| pb_id == filter_id)
});
Self { id, events }
}
pub fn id(&self) -> &str {
&self.id
}
pub fn into_events(self) -> FilteredSubscription<AriMessage> {
self.events
}
}