dahua_camera_rtsp/
camera_service.rs1use crate::frame_router::{FrameRouter, Subscription};
4use crate::ports::VideoSourceFactory;
5use dahua_camera_core::ports::{CameraControl, CameraControlFactory, EventSourceFactory};
6use crate::supervisor::{supervise, supervise_link, RestartPolicy, SupervisedLink, SupervisorCtx};
7use dahua_camera_core::{
8 AppConfig, CameraConfig, CameraStats, CameraStatus, CapabilitySet, Event, SharedParams,
9};
10use std::collections::BTreeMap;
11use std::sync::Arc;
12use std::time::Duration;
13use tokio::sync::broadcast;
14use tokio_util::sync::CancellationToken;
15
16const EVENT_BUS_CAPACITY: usize = 128;
18
19
20#[derive(Clone)]
26pub struct ServiceDeps {
27 pub source_factory: Arc<dyn VideoSourceFactory>,
29 pub control_factory: Arc<dyn CameraControlFactory>,
31 pub event_factory: Option<Arc<dyn EventSourceFactory>>,
33 #[cfg(feature = "decode")]
35 pub codec_factory: Arc<dyn crate::ports::CodecFactory>,
36 pub policy: RestartPolicy,
38}
39
40pub struct Camera {
45 id: String,
46 config: CameraConfig,
47 stats: Arc<CameraStats>,
48 router: Arc<FrameRouter>,
49 control: Arc<dyn CameraControl>,
50 capabilities: std::sync::RwLock<CapabilitySet>,
56 #[cfg(feature = "decode")]
57 jpeg: tokio::sync::Mutex<Option<crate::decoder::JpegLane>>,
58 #[cfg(feature = "decode")]
59 codec_factory: Arc<dyn crate::ports::CodecFactory>,
60}
61
62impl Camera {
63 pub fn id(&self) -> &str {
65 &self.id
66 }
67
68 pub fn config(&self) -> &CameraConfig {
70 &self.config
71 }
72
73 pub fn stats(&self) -> &CameraStats {
75 &self.stats
76 }
77
78 pub fn status(&self) -> CameraStatus {
80 self.router.status()
81 }
82
83 pub fn params(&self) -> Option<SharedParams> {
85 self.router.params()
86 }
87
88 pub fn viewer_count(&self) -> usize {
90 self.router.viewer_count()
91 }
92
93 pub fn control(&self) -> &Arc<dyn CameraControl> {
95 &self.control
96 }
97
98 pub fn capabilities(&self) -> CapabilitySet {
104 self.capabilities.read().expect("capabilities lock").clone()
105 }
106
107 pub fn set_capabilities(&self, capabilities: CapabilitySet) {
109 *self.capabilities.write().expect("capabilities lock") = capabilities;
110 }
111
112 pub fn subscribe(&self) -> Subscription {
118 self.router.subscribe(self.stats.clone())
119 }
120
121 pub async fn wait_for_params(&self, timeout: Duration) -> Option<SharedParams> {
126 self.router.wait_for_params(timeout).await
127 }
128}
129
130#[cfg(feature = "decode")]
131impl Camera {
132 pub async fn jpeg_lane(&self) -> crate::decoder::JpegLane {
138 let mut slot = self.jpeg.lock().await;
139 if let Some(lane) = slot.as_ref() {
140 if lane.is_alive() {
141 return lane.clone();
142 }
143 }
144 let lane = crate::decoder::spawn_lane(
145 &self.id,
146 self.config.jpeg_quality,
147 self.subscribe(),
148 self.stats.clone(),
149 self.codec_factory.clone(),
150 );
151 *slot = Some(lane.clone());
152 lane
153 }
154}
155
156pub struct CameraService {
158 cameras: BTreeMap<String, Arc<Camera>>,
159 events_tx: broadcast::Sender<Event>,
160 cancel: CancellationToken,
161 spawned_tasks: std::sync::atomic::AtomicUsize,
168}
169
170impl CameraService {
171 pub fn start(config: &AppConfig, deps: ServiceDeps) -> Arc<Self> {
176 let cancel = CancellationToken::new();
177 let (events_tx, _) = broadcast::channel(EVENT_BUS_CAPACITY);
178 let mut cameras = BTreeMap::new();
179 let mut spawned = 0usize;
180
181 for (id, camera_config) in &config.cameras {
182 let router = Arc::new(FrameRouter::new());
183 let stats = Arc::new(CameraStats::default());
184
185 cameras.insert(
186 id.clone(),
187 Arc::new(Camera {
188 id: id.clone(),
189 config: camera_config.clone(),
190 stats: stats.clone(),
191 router: router.clone(),
192 control: deps.control_factory.create(camera_config),
193 capabilities: std::sync::RwLock::new(CapabilitySet::default()),
194 #[cfg(feature = "decode")]
195 jpeg: tokio::sync::Mutex::new(None),
196 #[cfg(feature = "decode")]
197 codec_factory: deps.codec_factory.clone(),
198 }),
199 );
200
201 spawned += 1;
202 tokio::spawn(supervise(SupervisorCtx {
203 id: id.clone(),
204 config: camera_config.clone(),
205 router,
206 stats,
207 source_factory: deps.source_factory.clone(),
208 cancel: cancel.child_token(),
209 policy: deps.policy.clone(),
210 }));
211
212 if let Some(factory) = &deps.event_factory {
216 let link = Arc::new(EventLink {
217 source: factory.create(camera_config),
218 sink: events_tx.clone(),
219 });
220 spawned += 1;
221 tokio::spawn(supervise_link(
222 id.clone(),
223 link,
224 cancel.child_token(),
225 deps.policy.clone(),
226 ));
227 }
228 }
229
230 Arc::new(Self {
231 cameras,
232 events_tx,
233 cancel,
234 spawned_tasks: std::sync::atomic::AtomicUsize::new(spawned),
235 })
236 }
237
238 pub fn get(&self, id: &str) -> Option<Arc<Camera>> {
240 self.cameras.get(id).cloned()
241 }
242
243 pub fn all(&self) -> impl Iterator<Item = &Arc<Camera>> {
245 self.cameras.values()
246 }
247
248 pub fn subscribe_events(&self) -> broadcast::Receiver<Event> {
250 self.events_tx.subscribe()
251 }
252
253 pub fn spawned_task_count(&self) -> usize {
259 self.spawned_tasks.load(std::sync::atomic::Ordering::Relaxed)
260 }
261
262 pub fn shutdown(&self) {
264 self.cancel.cancel();
265 }
266}
267
268struct EventLink {
277 source: Arc<dyn dahua_camera_core::ports::EventSource>,
278 sink: broadcast::Sender<Event>,
279}
280
281#[async_trait::async_trait]
282impl SupervisedLink for EventLink {
283 fn name(&self) -> &'static str {
284 "events"
285 }
286
287 async fn run_once(&self, cancel: CancellationToken) -> dahua_camera_core::Result<()> {
288 let mut stream = self.source.connect().await?;
289 tracing::info!("event stream connected");
290
291 loop {
292 let next = tokio::select! {
293 biased;
294 _ = cancel.cancelled() => return Ok(()),
295 r = stream.next() => r,
296 };
297
298 match next? {
299 Some(event) => {
303 let _ = self.sink.send(event);
304 }
305 None => {
306 tracing::debug!("event stream closed by peer");
307 return Err(dahua_camera_core::CameraError::Cgi {
308 camera_id: String::new(),
309 detail: "event stream closed".to_owned(),
310 });
311 }
312 }
313 }
314 }
315}