Skip to main content

braid_run/
multicam_http_session_handler.rs

1// Copyright (C) The Strand-Braid Authors
2// SPDX-License-Identifier: MIT OR Apache-2.0
3
4use preferences_serde1::Preferences;
5use std::{
6    collections::BTreeMap,
7    sync::{Arc, RwLock},
8};
9use tracing::{debug, error, info, warn};
10
11use braid_types::{BuiServerInfo, RawCamName};
12use strand_bui_backend_session::HttpSession;
13use strand_cam_storetype::CallbackType;
14
15/// Keeps HTTP sessions for all connected cameras.
16#[derive(Clone)]
17pub(crate) struct StrandCamHttpSessionHandler {
18    cam_manager: flydra2::ConnectedCamerasManager,
19    pub(crate) name_to_session: Arc<RwLock<BTreeMap<RawCamName, MaybeSession>>>,
20    jar: Arc<RwLock<cookie_store::CookieStore>>,
21}
22
23#[derive(Clone)]
24pub(crate) enum MaybeSession {
25    Alive(HttpSession),
26    Errored,
27}
28
29use crate::mainbrain::{MainbrainError, MainbrainResult};
30
31impl StrandCamHttpSessionHandler {
32    pub(crate) fn new(
33        cam_manager: flydra2::ConnectedCamerasManager,
34        jar: Arc<RwLock<cookie_store::CookieStore>>,
35    ) -> Self {
36        Self {
37            cam_manager,
38            name_to_session: Arc::new(RwLock::new(BTreeMap::new())),
39            jar,
40        }
41    }
42    async fn open_session(&self, cam_name: &RawCamName) -> Result<MaybeSession, MainbrainError> {
43        // Create a new session if it doesn't exist.
44        let bui_server_addr_info = {
45            if let Some(cam_addr) = self.cam_manager.http_camserver_info(cam_name) {
46                match cam_addr {
47                    BuiServerInfo::NoServer => {
48                        panic!("cannot connect to camera with no server");
49                    }
50                    BuiServerInfo::Server(details) => details,
51                }
52            } else {
53                return Err(MainbrainError::UnknownCamera {
54                    cam_name: cam_name.clone(),
55                });
56            }
57        };
58
59        info!(
60            "opening session for cam {} to {}",
61            cam_name.as_str(),
62            bui_server_addr_info.addr(),
63        );
64
65        let result =
66            strand_bui_backend_session::create_session(&bui_server_addr_info, self.jar.clone())
67                .await;
68        let session = match result {
69            Ok(session) => {
70                let mut name_to_session = self.name_to_session.write().unwrap();
71                let session = MaybeSession::Alive(session);
72                name_to_session.insert(cam_name.clone(), session.clone());
73                session
74            }
75            Err(e) => {
76                error!(
77                    "could not create session to {}: {}",
78                    bui_server_addr_info.addr(),
79                    e
80                );
81                return Err(e.into());
82            }
83        };
84        {
85            // We have the cookie from braid now, so store it to disk.
86            let jar = self.jar.read().unwrap();
87            Preferences::save(
88                &*jar,
89                &crate::mainbrain::APP_INFO,
90                crate::mainbrain::STRAND_CAM_COOKIE_KEY,
91            )?;
92            // The jar holds live session cookies; keep its file owner-only.
93            braid_types::harden_prefs_file(
94                &crate::mainbrain::APP_INFO,
95                crate::mainbrain::STRAND_CAM_COOKIE_KEY,
96            );
97            tracing::debug!(
98                "saved cookie store {}",
99                crate::mainbrain::STRAND_CAM_COOKIE_KEY
100            );
101        }
102        Ok(session)
103    }
104
105    pub(crate) async fn get_or_open_session(
106        &self,
107        cam_name: &RawCamName,
108    ) -> Result<MaybeSession, MainbrainError> {
109        // Get session if it already exists.
110        let opt_session = { self.name_to_session.read().unwrap().get(cam_name).cloned() };
111
112        // Create session if needed.
113        match opt_session {
114            Some(session) => Ok(session),
115            None => self.open_session(cam_name).await,
116        }
117    }
118
119    async fn post(
120        &self,
121        cam_name: &RawCamName,
122        args: strand_cam_remote_control::CamArg,
123    ) -> Result<(), MainbrainError> {
124        self.post_callback(cam_name, CallbackType::ToCamera(args))
125            .await
126    }
127
128    async fn post_callback(
129        &self,
130        cam_name: &RawCamName,
131        callback: CallbackType,
132    ) -> Result<(), MainbrainError> {
133        let session = self.get_or_open_session(cam_name).await?;
134
135        // Post to session
136        match session {
137            MaybeSession::Alive(mut session) => {
138                let body = axum::body::Body::new(http_body_util::Full::new(bytes::Bytes::from(
139                    serde_json::to_vec(&callback).unwrap(),
140                )));
141
142                let result = session.post("callback", body).await;
143                match result {
144                    Ok(response) => {
145                        debug!(
146                            "StrandCamHttpSessionHandler::post() got response {:?}",
147                            response
148                        );
149                    }
150                    Err(err) => {
151                        error!(
152                            "For \"{}\": StrandCamHttpSessionHandler::post() got error {err:?}",
153                            cam_name.as_str(),
154                        );
155                        let mut name_to_session = self.name_to_session.write().unwrap();
156                        name_to_session.insert(cam_name.clone(), MaybeSession::Errored);
157                        // return Err(MainbrainError::blarg);
158                    }
159                }
160            }
161            MaybeSession::Errored => {
162                // TODO: should an error be raised here?
163                // return Err(MainbrainError::blarg);
164            }
165        };
166        Ok(())
167    }
168
169    pub(crate) async fn send_frame_offset(
170        &self,
171        cam_name: &RawCamName,
172        frame_offset: u64,
173    ) -> Result<(), MainbrainError> {
174        info!(
175            "for cam {}, sending frame offset {}",
176            cam_name.as_str(),
177            frame_offset
178        );
179        let args = strand_cam_remote_control::CamArg::SetFrameOffset(frame_offset);
180        self.post(cam_name, args).await
181    }
182
183    async fn send_quit(&mut self, cam_name: &RawCamName) -> Result<(), MainbrainError> {
184        info!("for cam {}, sending quit", cam_name.as_str());
185        let args = strand_cam_remote_control::CamArg::DoQuit;
186
187        let cam_result = self.post(cam_name, args).await;
188
189        // If we are telling the camera to quit, we don't want to keep its session around
190        let mut name_to_session = self.name_to_session.write().unwrap();
191        name_to_session.remove(cam_name);
192        self.cam_manager.remove(cam_name);
193        // TODO: we should cancel the stream of incoming frames so that they
194        // don't get processed after we have removed this camera
195        // information.
196
197        match cam_result {
198            Ok(_) => Ok(()),
199            Err(e) => {
200                warn!(
201                    "Ignoring error while sending quit command to \"{}\": {}",
202                    cam_name.as_str(),
203                    e
204                );
205                Err(e)
206            }
207        }
208    }
209
210    pub(crate) async fn send_quit_all(&mut self) {
211        use futures::{StreamExt, stream};
212        // Based on https://stackoverflow.com/a/51047786
213        const CONCURRENT_REQUESTS: usize = 5;
214        let results = stream::iter(self.cam_manager.all_raw_cam_names())
215            .map(|cam_name| {
216                let mut session = self.clone();
217                let cam_name = cam_name.clone();
218                async move {
219                    session
220                        .send_quit(&cam_name)
221                        .await
222                        .map_err(|e| (cam_name, e))
223                }
224            })
225            .buffer_unordered(CONCURRENT_REQUESTS);
226
227        results
228            .for_each(|r| async {
229                match r {
230                    Ok(()) => {}
231                    Err((cam_name, e)) => warn!(
232                        "Ignoring error When sending quit command to camera \"{}\": {}",
233                        cam_name.as_str(),
234                        e
235                    ),
236                }
237            })
238            .await;
239    }
240
241    pub(crate) async fn toggle_saving_mp4_files_all(
242        &self,
243        start_saving: bool,
244    ) -> MainbrainResult<()> {
245        let cam_names = self.cam_manager.all_raw_cam_names();
246        for cam_name in cam_names.iter() {
247            self.toggle_saving_mp4_files(cam_name, start_saving).await?;
248        }
249        Ok(())
250    }
251
252    pub(crate) async fn toggle_saving_mp4_files(
253        &self,
254        cam_name: &RawCamName,
255        start_saving: bool,
256    ) -> MainbrainResult<()> {
257        debug!(
258            "for cam {}, sending save mp4 file {:?}",
259            cam_name.as_str(),
260            start_saving
261        );
262        let cam_name = cam_name.clone();
263
264        let args = strand_cam_remote_control::CamArg::SetIsRecordingMp4(start_saving);
265        self.post(&cam_name, args).await?;
266        Ok(())
267    }
268
269    pub(crate) async fn send_clock_model_to_all(
270        &self,
271        clock_model: Option<strand_cam_bui_types::ClockModel>,
272    ) -> MainbrainResult<()> {
273        let cam_names = self.cam_manager.all_raw_cam_names();
274        for cam_name in cam_names.iter() {
275            self.send_triggerbox_clock_model(cam_name, clock_model.clone())
276                .await?;
277        }
278        Ok(())
279    }
280
281    pub(crate) async fn send_triggerbox_clock_model(
282        &self,
283        cam_name: &RawCamName,
284        clock_model: Option<strand_cam_bui_types::ClockModel>,
285    ) -> MainbrainResult<()> {
286        debug!(
287            "for cam {}, sending clock model {:?}",
288            cam_name.as_str(),
289            clock_model
290        );
291        let cam_name = cam_name.clone();
292
293        let args = strand_cam_remote_control::CamArg::SetTriggerboxClockModel(clock_model);
294        self.post(&cam_name, args).await
295    }
296
297    pub(crate) async fn set_post_trigger_buffer_all(
298        &self,
299        num_frames: usize,
300    ) -> MainbrainResult<()> {
301        let cam_names = self.cam_manager.all_raw_cam_names();
302        for cam_name in cam_names.iter() {
303            self.set_post_trigger_buffer(cam_name, num_frames).await?;
304        }
305        Ok(())
306    }
307
308    pub(crate) async fn set_post_trigger_buffer(
309        &self,
310        cam_name: &RawCamName,
311        num_frames: usize,
312    ) -> MainbrainResult<()> {
313        debug!(
314            "for cam {}, sending set post trigger buffer {}",
315            cam_name.as_str(),
316            num_frames
317        );
318        let cam_name = cam_name.clone();
319
320        let args = strand_cam_remote_control::CamArg::SetPostTriggerBufferSize(num_frames);
321        self.post(&cam_name, args).await?;
322        Ok(())
323    }
324
325    pub(crate) async fn initiate_post_trigger_mp4_all(&self) -> MainbrainResult<()> {
326        let cam_names = self.cam_manager.all_raw_cam_names();
327        for cam_name in cam_names.iter() {
328            self.initiate_post_trigger_mp4(cam_name).await?;
329        }
330        Ok(())
331    }
332
333    pub(crate) async fn initiate_post_trigger_mp4(
334        &self,
335        cam_name: &RawCamName,
336    ) -> MainbrainResult<()> {
337        debug!(
338            "for cam {}, initiating post trigger recording",
339            cam_name.as_str(),
340        );
341        let cam_name = cam_name.clone();
342
343        let args = strand_cam_remote_control::CamArg::PostTrigger;
344        self.post(&cam_name, args).await?;
345        Ok(())
346    }
347
348    pub(crate) async fn take_new_background_all(&self) -> MainbrainResult<()> {
349        let cam_names = self.cam_manager.all_raw_cam_names();
350        for cam_name in cam_names.iter() {
351            self.take_new_background(cam_name).await?;
352        }
353        Ok(())
354    }
355
356    pub(crate) async fn take_new_background(&self, cam_name: &RawCamName) -> MainbrainResult<()> {
357        debug!("for cam {}, taking new background image", cam_name.as_str());
358        self.post_callback(cam_name, CallbackType::TakeCurrentImageAsBackground)
359            .await?;
360        Ok(())
361    }
362
363    pub(crate) async fn send_obj_detection_config(
364        &self,
365        cam_name: &RawCamName,
366        cfg_yaml: String,
367    ) -> MainbrainResult<()> {
368        debug!(
369            "for cam {}, sending object detection config",
370            cam_name.as_str()
371        );
372        let args = strand_cam_remote_control::CamArg::SetObjDetectionConfig(cfg_yaml);
373        self.post(cam_name, args).await?;
374        Ok(())
375    }
376}