arcly_stream/bus/
application.rs1use super::events::{StreamEvent, StreamEventKind};
4use super::handle::{StreamHandle, StreamState};
5use crate::observe::Observer;
6use crate::{AppName, Result, StreamError, StreamId};
7use dashmap::DashMap;
8use std::sync::atomic::{AtomicUsize, Ordering};
9use std::sync::Arc;
10use tokio::sync::broadcast;
11use tracing::info;
12
13pub struct Application {
19 pub name: AppName,
21 pub broadcast_capacity: usize,
23 pub gop_capacity: usize,
25 pub gop_byte_capacity: usize,
27 pub subscriber_max_lag: Option<u64>,
29 streams: DashMap<StreamId, StreamHandle>,
30 event_tx: broadcast::Sender<StreamEvent>,
32 pub_lock: tokio::sync::Mutex<()>,
35 stream_count: AtomicUsize,
37 observer: Arc<dyn Observer>,
39}
40
41impl Application {
42 pub fn new(
44 name: AppName,
45 broadcast_capacity: usize,
46 gop_capacity: usize,
47 gop_byte_capacity: usize,
48 subscriber_max_lag: Option<u64>,
49 observer: Arc<dyn Observer>,
50 ) -> Arc<Self> {
51 let (event_tx, _) = broadcast::channel(64);
52 Arc::new(Self {
53 name,
54 broadcast_capacity,
55 gop_capacity,
56 gop_byte_capacity,
57 subscriber_max_lag,
58 streams: DashMap::new(),
59 event_tx,
60 pub_lock: tokio::sync::Mutex::new(()),
61 stream_count: AtomicUsize::new(0),
62 observer,
63 })
64 }
65
66 pub async fn start_publish(&self, stream_id: StreamId) -> Result<StreamHandle> {
69 let _guard = self.pub_lock.lock().await;
72
73 if let Some(existing) = self.streams.get(&stream_id) {
74 let state = existing.current_state().await;
75 if matches!(state, StreamState::Publishing | StreamState::Transcoding) {
76 return Err(StreamError::StreamAlreadyPublishing {
77 app: self.name.to_string(),
78 stream_id: stream_id.to_string(),
79 });
80 }
81 }
82
83 let handle = StreamHandle::with_config(
84 self.name.clone(),
85 stream_id.clone(),
86 self.broadcast_capacity,
87 self.gop_capacity,
88 self.gop_byte_capacity,
89 self.subscriber_max_lag,
90 Arc::clone(&self.observer),
91 );
92 handle.set_state(StreamState::Publishing).await;
93 let started_at_ms = super::handle::now_ms();
94 handle
95 .update_metadata(|m| m.started_at_ms = started_at_ms)
96 .await;
97 self.streams.insert(stream_id.clone(), handle.clone());
98 self.stream_count.fetch_add(1, Ordering::Relaxed);
99 drop(_guard);
102 self.observer.on_publish_started(self.name.as_str());
103
104 info!(app = %self.name, stream = %stream_id, "Stream publish started");
105 self.emit(stream_id.clone(), StreamEventKind::PublishStarted);
106
107 Ok(handle)
108 }
109
110 pub async fn end_publish(&self, stream_id: &StreamId) -> Result<bool> {
114 if let Some((_, handle)) = self.streams.remove(stream_id) {
115 self.stream_count.fetch_sub(1, Ordering::Relaxed);
116 self.observer.on_publish_ended(self.name.as_str());
117 handle.set_state(StreamState::Ended).await;
118 handle.close();
122 info!(app = %self.name, stream = %stream_id, "Stream publish ended");
123 self.emit(stream_id.clone(), StreamEventKind::PublishEnded);
124 return Ok(true);
125 }
126 Ok(false)
127 }
128
129 pub fn get_stream(&self, stream_id: &StreamId) -> Option<StreamHandle> {
131 self.streams.get(stream_id).map(|r| r.clone())
132 }
133
134 pub fn active_streams(&self) -> Vec<StreamId> {
136 self.streams.iter().map(|r| r.key().clone()).collect()
137 }
138
139 pub fn active_handles(&self) -> Vec<StreamHandle> {
141 self.streams.iter().map(|r| r.value().clone()).collect()
142 }
143
144 pub fn subscribe_events(&self) -> broadcast::Receiver<StreamEvent> {
146 self.event_tx.subscribe()
147 }
148
149 pub fn stream_count(&self) -> usize {
151 self.stream_count.load(Ordering::Relaxed)
152 }
153
154 fn emit(&self, stream_id: StreamId, kind: StreamEventKind) {
156 let event = StreamEvent {
157 app: self.name.clone(),
158 stream_id,
159 kind,
160 };
161 self.observer.on_event(&event);
162 let _ = self.event_tx.send(event);
163 }
164}
165
166impl std::fmt::Debug for Application {
167 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
168 f.debug_struct("Application")
169 .field("name", &self.name)
170 .field("stream_count", &self.streams.len())
171 .finish()
172 }
173}