Skip to main content

braid_run/
mainbrain.rs

1// Copyright (C) The Strand-Braid Authors
2// SPDX-License-Identifier: MIT OR Apache-2.0
3
4use std::{
5    collections::BTreeMap,
6    net::SocketAddr,
7    path::PathBuf,
8    sync::{
9        Arc, RwLock,
10        atomic::{AtomicBool, Ordering},
11    },
12};
13
14use async_change_tracker::ChangeTracker;
15use axum::{
16    extract::{Path, State},
17    routing::get,
18};
19use futures::StreamExt;
20use http::{HeaderValue, StatusCode};
21use preferences_serde1::{AppInfo, Preferences};
22use tokio::net::UdpSocket;
23use tower_http::trace::TraceLayer;
24use tracing::{debug, error, info};
25
26use braid_types::{
27    BRAID_EVENT_NAME, BRAID_EVENTS_URL_PATH, BRAID_QUIT_EVENT_NAME, BraidHttpApiSharedState,
28    CamInfo, CborPacketCodec, FakeSyncConfig, FlydraFloatTimestampLocal, PerCamSaveData,
29    RawCamName, SyncFno, TRIGGERBOX_SYNC_SECONDS, TriggerType, Triggerbox,
30    braid_http::{CAM_PROXY_PATH, REMOTE_CAMERA_INFO_PATH},
31};
32use event_stream_types::{AcceptsEventStream, EventBroadcaster};
33use flydra2::{CoordProcessor, CoordProcessorConfig, FrameDataAndPoints, StreamItem};
34use strand_bui_backend_session_types::AccessToken;
35use strand_bui_backend_session_types::BuiServerAddrInfo;
36use strand_cam_bui_types::{ClockModel, RecordingPath};
37
38use eyre::{self, Result, WrapErr};
39
40use crate::multicam_http_session_handler::{MaybeSession, StrandCamHttpSessionHandler};
41
42#[cfg(feature = "bundle_files")]
43static ASSETS_DIR: include_dir::Dir<'static> =
44    include_dir::include_dir!("$CARGO_MANIFEST_DIR/braid_frontend/dist");
45
46lazy_static::lazy_static! {
47    static ref EVENTS_PREFIX: String = format!("/{}", BRAID_EVENTS_URL_PATH);
48}
49
50pub(crate) const APP_INFO: AppInfo = AppInfo {
51    name: "braid",
52    author: "AndrewStraw",
53};
54const COOKIE_SECRET_KEY: &str = "cookie-secret-base64";
55pub(crate) const STRAND_CAM_COOKIE_KEY: &str = "strand-cam-cookie";
56
57type SharedStore = Arc<RwLock<ChangeTracker<BraidHttpApiSharedState>>>;
58
59#[derive(thiserror::Error, Debug)]
60pub(crate) enum MainbrainError {
61    #[error("{source}")]
62    HyperError {
63        #[from]
64        source: hyper::Error,
65    },
66    #[error("{source}")]
67    BuiBackendSessionError {
68        #[from]
69        source: strand_bui_backend_session::Error,
70    },
71    #[error("{source}")]
72    PreferencesError {
73        #[from]
74        source: preferences_serde1::PreferencesError,
75    },
76    #[error("unknown camera \"{cam_name}\"")]
77    UnknownCamera { cam_name: RawCamName },
78}
79
80pub(crate) type MainbrainResult<T> = std::result::Result<T, MainbrainError>;
81
82/// The structure that holds our app data
83#[derive(Clone)]
84pub(crate) struct BraidAppState {
85    pub(crate) shared_store: SharedStore,
86    lowlatency_camdata_udp_addr: SocketAddr,
87    force_camera_sync_mode: bool,
88    software_limit_framerate: braid_types::StartSoftwareFrameRateLimit,
89    event_broadcaster: EventBroadcaster<usize>,
90    pub(crate) per_cam_data_arc: Arc<RwLock<BTreeMap<RawCamName, PerCamSaveData>>>,
91    pub(crate) expected_framerate_arc: Arc<RwLock<Option<f32>>>,
92    camera_configs: BTreeMap<RawCamName, braid_types::BraidCameraConfig>,
93    next_connection_id: Arc<RwLock<usize>>,
94    pub(crate) strand_cam_http_session_handler: StrandCamHttpSessionHandler,
95    pub(crate) cam_manager: flydra2::ConnectedCamerasManager,
96    pub(crate) output_base_dirname: PathBuf,
97    pub(crate) braidz_write_tx_weak: tokio::sync::mpsc::WeakSender<flydra2::SaveToDiskMsg>,
98    /// Sending `()` initiates the graceful shutdown sequence.
99    pub(crate) shtdwn_q_tx: tokio::sync::mpsc::Sender<()>,
100    /// The address the HTTP server is bound to (possibly unspecified, e.g.
101    /// `0.0.0.0`), used to enumerate device-connection URLs on demand.
102    pub(crate) bui_server_info: BuiServerAddrInfo,
103    /// The cookie/token secret, used to mint a fresh short-lived access token
104    /// when a device-connection QR code is requested.
105    pub(crate) persistent_secret: cookie::Key,
106}
107
108async fn events_handler(
109    State(app_state): State<BraidAppState>,
110    session_key: axum_token_auth::SessionKey,
111    _: AcceptsEventStream,
112) -> impl axum::response::IntoResponse {
113    session_key.is_present();
114    let key = {
115        let mut next_connection_id = app_state.next_connection_id.write().unwrap();
116        let key = *next_connection_id;
117        *next_connection_id += 1;
118        key
119    };
120    let (tx, body) = app_state.event_broadcaster.new_connection(key);
121
122    // Send an initial copy of our state.
123    {
124        let current_state = app_state.shared_store.read().unwrap().as_ref().clone();
125        let frame_string = to_event_frame(&current_state);
126        match tx.send(http_body::Frame::data(frame_string.into())).await {
127            Ok(()) => {}
128            Err(_) => {
129                // The receiver was dropped because the connection closed. Should probably do more here.
130                tracing::debug!("initial send error");
131            }
132        }
133    }
134
135    body
136}
137
138async fn handle_auth_error(err: tower::BoxError) -> (StatusCode, &'static str) {
139    match err.downcast::<axum_token_auth::ValidationErrors>() {
140        Ok(err) => {
141            tracing::error!(
142                "Validation error(s): {:?}",
143                err.errors().collect::<Vec<_>>()
144            );
145            (StatusCode::UNAUTHORIZED, "Request is not authorized")
146        }
147        Err(orig_err) => {
148            tracing::error!("Unhandled internal error: {orig_err}");
149            (StatusCode::INTERNAL_SERVER_ERROR, "internal server error")
150        }
151    }
152}
153
154/// Query the mainbrain configuration to get data required for camera settings.
155///
156/// Note that this does not change the state of the mainbrain to register
157/// anything about the camera but only queries for its configuration.
158/// Registration of a new camera is done by
159/// [braid_types::BraidHttpApiCallback::NewCamera].
160async fn remote_camera_info_handler(
161    State(app_state): State<BraidAppState>,
162    session_key: axum_token_auth::SessionKey,
163    Path(raw_cam_name): Path<String>,
164) -> impl axum::response::IntoResponse {
165    session_key.is_present();
166    let cam_cfg = app_state
167        .camera_configs
168        .get(&RawCamName::new(raw_cam_name.clone()));
169
170    if let Some(config) = cam_cfg {
171        let software_limit_framerate = app_state.software_limit_framerate.clone();
172
173        let trig_config = app_state
174            .shared_store
175            .read()
176            .unwrap()
177            .as_ref()
178            .trigger_type
179            .clone();
180
181        let msg = braid_types::RemoteCameraInfoResponse {
182            camdata_udp_port: app_state.lowlatency_camdata_udp_addr.port(),
183            config: config.clone(),
184            force_camera_sync_mode: app_state.force_camera_sync_mode,
185            software_limit_framerate,
186            trig_config,
187        };
188        Ok(axum::Json(msg))
189    } else {
190        error!("HTTP camera not found: \"{raw_cam_name:?}\"");
191        Err((
192            StatusCode::NOT_FOUND,
193            format!("Camera \"{raw_cam_name}\" not found."),
194        ))
195    }
196}
197
198/// Returns the URLs (one per reachable network interface) at which this web UI
199/// can be reached, each carrying a freshly minted short-lived access token. The
200/// frontend turns these into QR codes for connecting another device.
201async fn device_connect_urls_handler(
202    State(app_state): State<BraidAppState>,
203    session_key: axum_token_auth::SessionKey,
204) -> impl axum::response::IntoResponse {
205    session_key.is_present();
206
207    let bound = *app_state.bui_server_info.addr();
208    // Match the token policy of `braid_types::start_listener`: a token is only
209    // required (and only useful) when the server is not bound to loopback.
210    let token = if bound.ip().is_loopback() {
211        AccessToken::NoToken
212    } else {
213        AccessToken::PreSharedToken(axum_token_auth::generate_token(
214            &app_state.persistent_secret,
215            braid_types::ACCESS_TOKEN_TTL,
216        ))
217    };
218    let info = BuiServerAddrInfo::new(bound, token);
219    let uris = match strand_bui_backend_session::build_urls(&info) {
220        Ok(uris) => uris,
221        Err(e) => {
222            return Err((
223                StatusCode::INTERNAL_SERVER_ERROR,
224                format!("failed to enumerate network interfaces: {e}"),
225            ));
226        }
227    };
228    let loopback_only = uris.iter().all(braid_types::is_loopback);
229    let urls = uris.into_iter().map(|u| u.to_string()).collect();
230    Ok(axum::Json(
231        strand_bui_backend_session_types::DeviceConnectUrls {
232            urls,
233            loopback_only,
234        },
235    ))
236}
237
238async fn cam_proxy_handler_inner(
239    app_state: BraidAppState,
240    session_key: axum_token_auth::SessionKey,
241    raw_cam_name: String,
242    cam_path: String,
243    req: axum::extract::Request,
244) -> impl axum::response::IntoResponse {
245    session_key.is_present();
246    tracing::debug!("raw_cam_name: {raw_cam_name}, cam_path: \"{cam_path}\", req: {req:?}");
247    let accepts: Vec<HeaderValue> = req
248        .headers()
249        .get_all(http::header::ACCEPT)
250        .iter()
251        .cloned()
252        .collect();
253    let cam_name = RawCamName::new(raw_cam_name);
254
255    let session = app_state
256        .strand_cam_http_session_handler
257        .get_or_open_session(&cam_name)
258        .await
259        .map_err(|e| match e {
260            MainbrainError::UnknownCamera { cam_name, .. } => {
261                let err_msg = format!("Unknown camera \"{cam_name}\"");
262                tracing::error!(err_msg);
263                (StatusCode::NOT_FOUND, err_msg)
264            }
265            _ => {
266                let err_msg = format!("Internal server error: {e} {e:?}");
267                tracing::error!(err_msg);
268                (StatusCode::INTERNAL_SERVER_ERROR, err_msg)
269            }
270        })?;
271
272    match session {
273        MaybeSession::Alive(mut session) => {
274            tracing::debug!("Will request path \"{cam_path}\". Got session {session:?}.");
275            session
276                .req_accepts(&cam_path, &accepts, req.method().clone(), req.into_body())
277                .await
278                .map_err(|e| {
279                    let err_msg = format!("Failed request to Strand Cam: {e} {e:?}");
280                    tracing::error!(err_msg);
281                    (StatusCode::INTERNAL_SERVER_ERROR, err_msg)
282                })
283        }
284        MaybeSession::Errored => Err((
285            StatusCode::INTERNAL_SERVER_ERROR,
286            format!(
287                "Braid lost connection to camera name \"{}\".",
288                cam_name.as_str()
289            ),
290        )),
291    }
292}
293
294async fn cam_proxy_handler_root(
295    State(app_state): State<BraidAppState>,
296    session_key: axum_token_auth::SessionKey,
297    Path(raw_cam_name): Path<String>,
298    req: axum::extract::Request,
299) -> impl axum::response::IntoResponse {
300    session_key.is_present();
301    cam_proxy_handler_inner(app_state, session_key, raw_cam_name, "".into(), req).await
302}
303
304async fn cam_proxy_handler(
305    State(app_state): State<BraidAppState>,
306    session_key: axum_token_auth::SessionKey,
307    Path((raw_cam_name, cam_path)): Path<(String, String)>,
308    req: axum::extract::Request,
309) -> impl axum::response::IntoResponse {
310    session_key.is_present();
311    cam_proxy_handler_inner(app_state, session_key, raw_cam_name, cam_path, req).await
312}
313
314/// Load the persistent cookie/token secret, generating and saving a fresh one
315/// if none exists.
316///
317/// This secret signs both the session cookies and the self-expiring access
318/// tokens, so it must be loaded once and shared between
319/// [`braid_types::start_listener`] (which mints Braid's token) and the auth
320/// layer that validates it. A stable secret across restarts is what keeps
321/// already-issued browser cookies — and Braid's persisted per-camera cookie
322/// jar — valid through an upgrade.
323pub(crate) fn load_persistent_secret(secret_override: Option<String>) -> Result<cookie::Key> {
324    use base64::Engine;
325    let persistent_secret_base64 = if let Some(secret) = secret_override {
326        secret
327    } else {
328        match String::load(&APP_INFO, COOKIE_SECRET_KEY) {
329            Ok(secret_base64) => secret_base64,
330            Err(_) => {
331                tracing::debug!("No secret loaded from preferences file, generating new.");
332                let persistent_secret = cookie::Key::generate();
333                let persistent_secret_base64 =
334                    base64::engine::general_purpose::STANDARD.encode(persistent_secret.master());
335                persistent_secret_base64.save(&APP_INFO, COOKIE_SECRET_KEY)?;
336                persistent_secret_base64
337            }
338        }
339    };
340
341    // The secret can forge any session cookie and mint any token, so ensure its
342    // on-disk file is owner-only.
343    braid_types::harden_prefs_file(&APP_INFO, COOKIE_SECRET_KEY);
344
345    let persistent_secret =
346        base64::engine::general_purpose::STANDARD.decode(persistent_secret_base64)?;
347    Ok(cookie::Key::try_from(persistent_secret.as_slice())?)
348}
349
350async fn launch_braid_http_backend(
351    persistent_secret: cookie::Key,
352    trusted_networks: Vec<axum_token_auth::CidrBlock>,
353    listener: tokio::net::TcpListener,
354    mainbrain_server_info: BuiServerAddrInfo,
355    app_state: BraidAppState,
356) -> Result<impl futures::Future<Output = Result<()>>> {
357    // Setup our auth layer. With self-expiring signed tokens the auth layer no
358    // longer stores a token value: it accepts any unexpired token signed with
359    // `persistent_secret`. We only need to know whether a token is required.
360    let token_config = match mainbrain_server_info.token() {
361        AccessToken::PreSharedToken(_) => Some(axum_token_auth::TokenConfig::new("token")),
362        AccessToken::NoToken => None,
363    };
364
365    // `AuthConfig` is `#[non_exhaustive]`, so build it via `new` and set fields.
366    let mut cfg = axum_token_auth::AuthConfig::new(persistent_secret);
367    cfg.token_config = token_config;
368    cfg.cookie_name = "braid-bui-session";
369    // Sessions slide forward on use and survive up to 400 days of absence,
370    // enforced server-side via the signed cookie. Existing cookies that
371    // predate this field carry no embedded expiry and are treated as
372    // non-expiring until renewed, so they stay valid across the upgrade.
373    cfg.session_expires = Some(std::time::Duration::from_secs(60 * 60 * 24 * 400)); // 400 days
374    // Clients on a trusted overlay network (e.g. Tailscale/WireGuard) are
375    // accepted without a token; the overlay has already authenticated them.
376    cfg.trusted_networks = trusted_networks;
377
378    #[cfg(feature = "bundle_files")]
379    let serve_dir = tower_serve_static::ServeDir::new(&ASSETS_DIR);
380
381    #[cfg(feature = "serve_files")]
382    let serve_dir = tower_http::services::fs::ServeDir::new(
383        std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
384            .join("braid_frontend")
385            .join("dist"),
386    );
387
388    let auth_layer = cfg.into_layer();
389
390    assert_eq!(BRAID_EVENTS_URL_PATH, "braid-events");
391    assert_eq!(REMOTE_CAMERA_INFO_PATH, "remote-camera-info");
392    assert_eq!(CAM_PROXY_PATH, "cam-proxy");
393
394    // Create axum router.
395    let router = axum::Router::new()
396        .route("/braid-events", get(events_handler))
397        .route(
398            "/remote-camera-info/{encoded_cam_name}",
399            get(remote_camera_info_handler),
400        )
401        .route("/device-connect-urls", get(device_connect_urls_handler))
402        // .route("/cam-proxy/:encoded_cam_name", get(slash_redirect_handler))
403        .route(
404            "/cam-proxy/{encoded_cam_name}/",
405            axum::routing::method_routing::any(cam_proxy_handler_root),
406        )
407        .route(
408            "/cam-proxy/{encoded_cam_name}/{*path}",
409            axum::routing::method_routing::any(cam_proxy_handler),
410        )
411        .route(
412            "/callback",
413            axum::routing::post(crate::callback_handling::callback_handler)
414                .layer(axum::extract::DefaultBodyLimit::max(100_000_000)),
415        )
416        .fallback_service(serve_dir)
417        .layer(
418            tower::ServiceBuilder::new()
419                .layer(TraceLayer::new_for_http())
420                // Auth layer will produce an error if the request cannot be
421                // authorized so we must handle that.
422                .layer(axum::error_handling::HandleErrorLayer::new(
423                    handle_auth_error,
424                ))
425                .layer(auth_layer),
426        )
427        .with_state(app_state);
428
429    // create future for our app
430    let http_serve_future = {
431        use futures::TryFutureExt;
432        use std::future::IntoFuture;
433        // `into_make_service_with_connect_info` exposes the peer address to the
434        // auth layer so it can recognize clients on a trusted overlay network.
435        axum::serve(
436            listener,
437            router.into_make_service_with_connect_info::<std::net::SocketAddr>(),
438        )
439        .into_future()
440        .map_err(eyre::Report::from)
441    };
442
443    // Display where we are listening.
444    info!(
445        "Braid HTTP server listening at {}",
446        mainbrain_server_info.addr()
447    );
448
449    let urls = strand_bui_backend_session::build_urls(&mainbrain_server_info)?;
450    for url in urls.iter() {
451        info!("Predicted URL: {url}");
452        if !braid_types::is_loopback(url) {
453            println!("QR code for {url}");
454            display_qr_url(&format!("{url}"))?;
455        }
456    }
457
458    Ok(http_serve_future)
459}
460
461fn compute_trigger_timestamp(
462    model: &Option<ClockModel>,
463    synced_frame: SyncFno,
464) -> Option<FlydraFloatTimestampLocal<Triggerbox>> {
465    if let Some(model) = model {
466        let v: f64 = (synced_frame.0 as f64) * model.gain + model.offset;
467        Some(FlydraFloatTimestampLocal::from_f64(v))
468    } else {
469        None
470    }
471}
472
473struct SendConnectedCamToBuiBackend {
474    shared_store: SharedStore,
475}
476
477impl flydra2::ConnectedCamCallback for SendConnectedCamToBuiBackend {
478    fn on_cam_changed(&self, new_cam_list: Vec<CamInfo>) {
479        let mut tracker = self.shared_store.write().unwrap();
480        tracker.modify(|shared| shared.connected_cameras = new_cam_list.clone());
481    }
482}
483
484fn display_qr_url(url: &str) -> Result<()> {
485    use qrcode::QrCode;
486    use qrcode::render::unicode;
487    use std::io::{Write, stdout};
488
489    let qr = QrCode::new(url)?;
490
491    let image = qr.render::<unicode::Dense1x2>().build();
492
493    let stdout = stdout();
494    let mut stdout_handle = stdout.lock();
495    writeln!(stdout_handle)?;
496    stdout_handle.write_all(image.as_bytes())?;
497    writeln!(stdout_handle)?;
498    Ok(())
499}
500
501#[allow(clippy::too_many_arguments)]
502pub(crate) async fn do_run_forever(
503    show_tracking_params: bool,
504    // sched_policy_priority: Option<(libc::c_int, libc::c_int)>,
505    camera_configs: BTreeMap<RawCamName, braid_types::BraidCameraConfig>,
506    trigger_cfg: TriggerType,
507    mainbrain_config: braid_config_data::MainbrainConfig,
508    persistent_secret: cookie::Key,
509    all_expected_cameras: std::collections::BTreeSet<RawCamName>,
510    force_camera_sync_mode: bool,
511    software_limit_framerate: braid_types::StartSoftwareFrameRateLimit,
512    saving_program_name: &str,
513    listener: tokio::net::TcpListener,
514    mainbrain_server_info: BuiServerAddrInfo,
515    mut strand_cam_set: tokio::task::JoinSet<()>,
516) -> Result<()> {
517    let cal_fname: Option<std::path::PathBuf> = mainbrain_config.cal_fname.clone();
518    let output_base_dirname: std::path::PathBuf = mainbrain_config.output_base_dirname.clone();
519    let tracking_params: braid_types::TrackingParams = mainbrain_config.tracking_params.clone();
520
521    let lowlatency_camdata_udp_port = &mainbrain_config.lowlatency_camdata_udp_port;
522    let mut ensure_camdata_ip = None;
523    if let Some(lowlatency_camdata_udp_addr) = &mainbrain_config.lowlatency_camdata_udp_addr {
524        tracing::warn!(
525            "Using deprecated configuration `lowlatency_camdata_udp_addr`. Use `lowlatency_camdata_udp_port` instead."
526        );
527        let lowlatency_camdata_udp_addr = lowlatency_camdata_udp_addr.parse::<SocketAddr>()?;
528        if lowlatency_camdata_udp_addr.port() != *lowlatency_camdata_udp_port {
529            eyre::bail!("camdata UDP port specified two different ways");
530        }
531        ensure_camdata_ip = Some(lowlatency_camdata_udp_addr.ip());
532    }
533
534    let save_empty_data2d: bool = mainbrain_config.save_empty_data2d;
535    let write_buffer_size_num_messages = mainbrain_config.write_buffer_size_num_messages;
536
537    info!("saving to directory: {}", output_base_dirname.display());
538
539    // Create `stream_cancel::Valve` for shutting everything down. Note this is
540    // `Clone`, so we can (and should) shut down everything with it.
541    let (quit_trigger, valve) = stream_cancel::Valve::new();
542    let (shtdwn_q_tx, mut shtdwn_q_rx) = tokio::sync::mpsc::channel::<()>(5);
543
544    // The SSE event broadcaster. Browser connections register here (via the app
545    // state) and state updates are pushed through it. It is created here, before
546    // the shutdown handler below, so that handler can broadcast a final "quit"
547    // event to every connected client during shutdown.
548    let event_broadcaster: EventBroadcaster<usize> = EventBroadcaster::default();
549
550    let recon = if let Some(ref cal_fname) = cal_fname {
551        info!("using calibration: {}", cal_fname.display());
552        let require_radfiles = false;
553        Some(
554            flydra_mvg::FlydraMultiCameraSystem::from_path(cal_fname, require_radfiles)
555                .with_context(|| {
556                    format!("loading calibration in file \"{}\"", cal_fname.display())
557                })?,
558        )
559    } else {
560        None
561    };
562
563    let signal_all_cams_present = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
564    let signal_all_cams_synced = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
565
566    let periodic_signal_period_usec = if let TriggerType::PtpSync(ptpcfg) = &trigger_cfg {
567        ptpcfg.periodic_signal_period_usec
568    } else {
569        None
570    };
571
572    let mut cam_manager = flydra2::ConnectedCamerasManager::new(
573        &recon,
574        all_expected_cameras,
575        signal_all_cams_present.clone(),
576        signal_all_cams_synced.clone(),
577        periodic_signal_period_usec,
578        None,
579    );
580
581    let jar: cookie_store::CookieStore = match Preferences::load(&APP_INFO, STRAND_CAM_COOKIE_KEY) {
582        Ok(jar) => {
583            tracing::debug!("loaded cookie store {STRAND_CAM_COOKIE_KEY}");
584            jar
585        }
586        Err(e) => {
587            tracing::debug!("cookie store {STRAND_CAM_COOKIE_KEY} not loaded: {e} {e:?}");
588            cookie_store::CookieStore::new(None)
589        }
590    };
591    let jar = Arc::new(RwLock::new(jar.clone()));
592    let strand_cam_http_session_handler =
593        StrandCamHttpSessionHandler::new(cam_manager.clone(), jar);
594
595    if show_tracking_params {
596        let t2: braid_types::TrackingParams = tracking_params;
597        let buf = toml::to_string(&t2)?;
598        println!("{}", buf);
599        std::process::exit(0);
600    }
601
602    let ignore_latency = false;
603    let mut coord_processor = CoordProcessor::new(
604        CoordProcessorConfig {
605            tracking_params,
606            save_empty_data2d,
607            ignore_latency,
608            mini_arena_debug_cfg: None,
609            write_buffer_size_num_messages,
610        },
611        cam_manager.clone(),
612        recon.clone(),
613        flydra2::BraidMetadataBuilder::saving_program_name(saving_program_name),
614    )?;
615
616    // Here is what we do on quit:
617    // 1) Stop saving data, convert .braid dir to .braidz, close files.
618    // 2) Fire a DoQuit message to all cameras and wait for them to quit.
619    // 3) Only then close all our network ports and streams nicely.
620    let mut quit_trigger_container = Some(quit_trigger);
621    let mut strand_cam_http_session_handler2 = strand_cam_http_session_handler.clone();
622    let braidz_write_tx_weak = coord_processor.braidz_write_tx.downgrade();
623    let shutdown_event_broadcaster = event_broadcaster.clone();
624    tokio::spawn(async move {
625        while let Some(()) = shtdwn_q_rx.recv().await {
626            debug!("got shutdown command {}:{}", file!(), line!());
627
628            if let Some(braidz_write_tx) = braidz_write_tx_weak.upgrade() {
629                // `braidz_write_tx` will be dropped after this scope.
630
631                // Stop saving Braid data.
632
633                // Do not need to wait for completion because we are going to
634                // exit nicely by manually ending all threads and letting all
635                // drop handlers run (without aborting) and thus the program
636                // will finish writing without an explicit wait. (Of course,
637                // this fails during an actual abort).
638
639                braidz_write_tx
640                    .send(flydra2::SaveToDiskMsg::StopSavingCsv)
641                    .await
642                    .unwrap_or(()); // ignore error on shutdown
643            }
644
645            strand_cam_http_session_handler2.send_quit_all().await;
646
647            // Tell every connected browser we are shutting down, so all clients
648            // (not only the one that pressed Quit) show the "Braid has quit"
649            // screen and stop reconnecting. This is sent while the HTTP server
650            // is still running: `quit_trigger.cancel()` below ends the
651            // top-level `select!` and thus drops the server future. The brief
652            // sleep gives the message time to flush to clients first.
653            shutdown_event_broadcaster
654                .broadcast_frame(quit_event_frame())
655                .await;
656            tokio::time::sleep(std::time::Duration::from_millis(500)).await;
657
658            // When we get here, we have successfully sent DoQuit to all cams.
659            // We can now quit everything in the mainbrain.
660            if let Some(quit_trigger) = quit_trigger_container.take() {
661                quit_trigger.cancel();
662                break; // no point to listen for more
663            }
664        }
665        debug!("shutdown handler finished {}:{}", file!(), line!());
666    });
667
668    let (triggerbox_cmd, triggerbox_rx) = match &trigger_cfg {
669        TriggerType::TriggerboxV1(_) => {
670            let (tx, rx) = tokio::sync::mpsc::channel(20);
671            (Some(tx), Some(rx))
672        }
673        TriggerType::FakeSync(_) | TriggerType::PtpSync(_) | TriggerType::DeviceTimestamp => {
674            (None, None)
675        }
676    };
677
678    let needs_clock_model = match &trigger_cfg {
679        TriggerType::TriggerboxV1(_) | TriggerType::FakeSync(_) => true,
680        TriggerType::PtpSync(_) | TriggerType::DeviceTimestamp => false,
681    };
682
683    let sync_pulse_pause_started: Option<flydra2::SyncStart> = None;
684    let sync_pulse_pause_started_arc = Arc::new(RwLock::new(sync_pulse_pause_started));
685
686    let flydra_app_name = "Braid".to_string();
687
688    let shared = BraidHttpApiSharedState {
689        trigger_type: trigger_cfg.clone(),
690        csv_tables_dirname: None,
691        fake_mp4_recording_path: None,
692        post_trigger_buffer_size: 0,
693        clock_model: None,
694        calibration_filename: cal_fname.map(|x| x.into_os_string().into_string().unwrap()),
695        connected_cameras: Vec::new(),
696        background_model_updating: Default::default(),
697        camera_image_dimensions: Default::default(),
698        model_server_addr: None,
699        flydra_app_name,
700        all_expected_cameras_are_synced: false,
701        needs_clock_model,
702        version_update: None,
703    };
704    let shared_store = ChangeTracker::new(shared);
705    let mut shared_store_changes_rx = shared_store.get_changes(1);
706    let shared_store = Arc::new(RwLock::new(shared_store));
707
708    // Periodically ask the version-check server whether a newer Braid release is
709    // available; if so, surface it to connected browsers as a dismissible
710    // banner. Disabled by setting DISABLE_VERSION_CHECK=1. Mirrors the check in
711    // Strand Camera.
712    let do_version_check = match std::env::var_os("DISABLE_VERSION_CHECK") {
713        Some(v) => &v == "0",
714        None => true,
715    };
716    if do_version_check {
717        let app_version: semver::Version = {
718            let mut my_version = semver::Version::parse(env!("CARGO_PKG_VERSION")).unwrap();
719            my_version.build = semver::BuildMetadata::new(env!("GIT_HASH"))?;
720            my_version
721        };
722        info!(
723            "This program will check for new versions automatically. To disable, \
724            set the environment variable DISABLE_VERSION_CHECK=1."
725        );
726        let store_for_version_check = shared_store.clone();
727        let checker = strand_version_check::VersionChecker::new();
728        let interval_stream = tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(
729            std::time::Duration::from_secs(1800),
730        ));
731        let mut interval_stream = valve.wrap(interval_stream);
732        tokio::spawn(async move {
733            // The newest version known so far; advanced as the server reports
734            // newer ones so each is announced only once.
735            let mut known_version = app_version;
736            while interval_stream.next().await.is_some() {
737                let user_agent = format!("braid/{}", known_version);
738                if let Some(av) = checker.fetch("braid", &user_agent).await
739                    && av.version > known_version
740                {
741                    info!(
742                        "New version of Braid is available: {}. {}",
743                        av.version, av.message
744                    );
745                    let update = braid_types::VersionUpdate {
746                        available: av.version.to_string(),
747                        message: av.message,
748                        url: av.url,
749                    };
750                    let mut tracker = store_for_version_check.write().unwrap();
751                    tracker.modify(|shared| shared.version_update = Some(update));
752                    known_version = av.version;
753                }
754            }
755        });
756    }
757
758    let expected_framerate_arc = Arc::new(RwLock::new(None));
759
760    let per_cam_data_arc = Arc::new(RwLock::new(Default::default()));
761
762    let (lowlatency_camdata_udp_addr, camdata_socket) = {
763        // The port of the low latency UDP incoming data socket may be specified
764        // as 0 in which case the OS will decide which port will actually be
765        // bound. So here we create the socket and get its port.
766        let camdata_addr_unspecified_port = {
767            // No low latency UDP port specified. Default to the same IP
768            // as the mainbrain HTTP server (which may be unspecified) and
769            // let the OS assign a free port by setting the port as
770            // unspecified.
771            let mainbrain_tcp_addr = listener.local_addr()?;
772            if let Some(ensure_camdata_ip) = ensure_camdata_ip
773                && mainbrain_tcp_addr.ip() != ensure_camdata_ip
774            {
775                eyre::bail!(
776                    "requested camdata UDP IP address not equal to mainbrain TCP IP address"
777                );
778            }
779            let mut camdata_addr_unspecified_port = mainbrain_tcp_addr;
780            camdata_addr_unspecified_port.set_port(*lowlatency_camdata_udp_port);
781            camdata_addr_unspecified_port
782        };
783        let camdata_socket = UdpSocket::bind(&camdata_addr_unspecified_port).await?;
784        let camdata_addr = camdata_socket.local_addr()?;
785        debug!("flydra mainbrain camera UDP listener socket: internal: {camdata_addr}");
786
787        (camdata_addr, camdata_socket)
788    };
789
790    if !output_base_dirname.exists() {
791        info!(
792            "creating output data directory at \"{}\"",
793            output_base_dirname.display()
794        );
795        std::fs::create_dir_all(&output_base_dirname)?;
796    }
797
798    debug!(
799        "output .braidz data directory will be \"{}\"",
800        std::fs::canonicalize(&output_base_dirname)?.display()
801    );
802
803    let braidz_write_tx_weak = coord_processor.braidz_write_tx.downgrade();
804
805    let time_model_arc = Arc::new(RwLock::new(None));
806
807    // Create our app state.
808    let app_state = BraidAppState {
809        shared_store: shared_store.clone(),
810        lowlatency_camdata_udp_addr,
811        force_camera_sync_mode,
812        software_limit_framerate,
813        event_broadcaster: event_broadcaster.clone(),
814        per_cam_data_arc: per_cam_data_arc.clone(),
815        camera_configs,
816        next_connection_id: Arc::new(RwLock::new(0)),
817        expected_framerate_arc: expected_framerate_arc.clone(),
818        braidz_write_tx_weak,
819        cam_manager: cam_manager.clone(),
820        output_base_dirname,
821        strand_cam_http_session_handler: strand_cam_http_session_handler.clone(),
822        shtdwn_q_tx,
823        bui_server_info: mainbrain_server_info.clone(),
824        persistent_secret: persistent_secret.clone(),
825    };
826
827    // This future will send state updates to all connected event listeners.
828    let event_broadcaster = app_state.event_broadcaster.clone();
829    let event_broadcast_fut = async move {
830        while let Some((_prev_state, next_state)) = shared_store_changes_rx.next().await {
831            let frame_string = to_event_frame(&next_state);
832            event_broadcaster.broadcast_frame(frame_string).await;
833        }
834    };
835
836    let trusted_networks = braid_types::parse_trusted_networks(&mainbrain_config.trusted_networks)?;
837    let http_serve_future = launch_braid_http_backend(
838        persistent_secret,
839        trusted_networks,
840        listener,
841        mainbrain_server_info,
842        app_state,
843    )
844    .await?;
845
846    let signal_triggerbox_connected = Arc::new(AtomicBool::new(false));
847
848    {
849        let sender = SendConnectedCamToBuiBackend {
850            shared_store: shared_store.clone(),
851        };
852        let old_callback = cam_manager.set_cam_changed_callback(Box::new(sender));
853        assert!(old_callback.is_none());
854    }
855
856    let (triggerbox_data_tx, mut triggerbox_data_rx) =
857        tokio::sync::mpsc::channel::<braid_triggerbox::TriggerClockInfoRow>(20);
858
859    match &trigger_cfg {
860        TriggerType::TriggerboxV1(_) | TriggerType::FakeSync(_) => {
861            let braidz_write_tx_weak = coord_processor.braidz_write_tx.downgrade();
862            let signal_triggerbox_connected = signal_triggerbox_connected.clone();
863
864            let mut has_triggerbox_connected = false;
865            let triggerbox_future = async move {
866                debug!(
867                    "starting triggerbox listener future {}:{}",
868                    file!(),
869                    line!()
870                );
871                while let Some(msg) = triggerbox_data_rx.recv().await {
872                    if !has_triggerbox_connected {
873                        has_triggerbox_connected = true;
874                        info!("triggerbox is connected.");
875                        signal_triggerbox_connected.store(true, Ordering::SeqCst);
876                    }
877                    let msg2 = braid_types::TriggerClockInfoRow {
878                        start_timestamp: msg.start_timestamp.into(),
879                        framecount: msg.framecount,
880                        tcnt: msg.tcnt,
881                        stop_timestamp: msg.stop_timestamp.into(),
882                    };
883
884                    if let Some(braidz_write_tx) = braidz_write_tx_weak.upgrade() {
885                        // `braidz_write_tx` will be dropped after this scope.
886                        braidz_write_tx
887                            .send(flydra2::SaveToDiskMsg::TriggerClockInfo(msg2))
888                            .await
889                            .unwrap();
890                    }
891                }
892                debug!("triggerbox listener future done {}:{}", file!(), line!());
893            };
894            tokio::spawn(triggerbox_future);
895        }
896        _ => {
897            debug!("not listening to triggerbox");
898        }
899    }
900
901    let tracker = shared_store.clone();
902
903    let on_new_clock_model = {
904        let time_model_arc = time_model_arc.clone();
905        let strand_cam_http_session_handler = strand_cam_http_session_handler.clone();
906        let tracker = tracker.clone();
907        let trigger_cfg = trigger_cfg.clone();
908        Box::new(move |tm1: Option<braid_triggerbox::ClockModel>| {
909            match &trigger_cfg {
910                TriggerType::FakeSync(_) | TriggerType::TriggerboxV1(_) => {
911                    let tm = tm1.map(|x| strand_cam_bui_types::ClockModel {
912                        gain: x.gain,
913                        offset: x.offset,
914                        n_measurements: x.n_measurements,
915                        residuals: x.residuals,
916                    });
917                    let cm = tm.clone();
918                    {
919                        let mut guard = time_model_arc.write().unwrap();
920                        *guard = tm;
921                    }
922                    {
923                        let mut tracker_guard = tracker.write().unwrap();
924                        tracker_guard.modify(|shared| shared.clock_model = cm.clone());
925                    }
926                    let strand_cam_http_session_handler2 = strand_cam_http_session_handler.clone();
927                    // TODO: Do we really need to spawn here? Why not just .await?
928                    tokio::spawn(async move {
929                        let r = strand_cam_http_session_handler2
930                            .send_clock_model_to_all(cm)
931                            .await;
932                        match r {
933                            Ok(_http_response) => {}
934                            Err(e) => {
935                                error!("error sending clock model: {}", e);
936                            }
937                        };
938                    });
939                }
940                TriggerType::PtpSync(_) | TriggerType::DeviceTimestamp => {
941                    // no central clock model
942                    panic!("No need for clock model.");
943                }
944            }
945        })
946    };
947
948    match &trigger_cfg {
949        TriggerType::TriggerboxV1(cfg) => {
950            let device_fname = cfg.device_fname.clone();
951            let fps = &cfg.framerate;
952            let query_dt = &cfg.query_dt;
953
954            use braid_triggerbox::{Cmd, make_trig_fps_cmd};
955
956            let tx = triggerbox_cmd.clone().unwrap();
957            let cmd_rx = triggerbox_rx.unwrap();
958
959            let (rate_cmd, rate_actual) = make_trig_fps_cmd(*fps as f64);
960
961            let max_triggerbox_measurement_error =
962                cfg.max_triggerbox_measurement_error.unwrap_or_else(|| {
963                    braid_types::TriggerboxConfig::default()
964                        .max_triggerbox_measurement_error
965                        .unwrap()
966                });
967
968            // queue several commands for the triggerbox on initial start.
969            tx.send(Cmd::StopPulsesAndReset).await?;
970            info!(
971                "Triggerbox at {} request {} fps, actual frame rate will be {} fps. Will \
972                accept maximum timestamp error of {} microseconds.",
973                device_fname,
974                fps,
975                rate_actual,
976                max_triggerbox_measurement_error.as_micros(),
977            );
978            tx.send(rate_cmd).await?;
979            tx.send(Cmd::StartPulses).await?;
980
981            {
982                let mut expected_framerate = expected_framerate_arc.write().unwrap();
983                *expected_framerate = Some(rate_actual as f32);
984            }
985
986            // Emperically, an Arduino Nano requires 7 seconds to wake up.
987            let sleep_dur = std::time::Duration::from_secs_f32(7.0);
988
989            let triggerbox = braid_triggerbox::TriggerboxDevice::new(
990                on_new_clock_model,
991                device_fname,
992                cmd_rx,
993                Some(triggerbox_data_tx),
994                None,
995                max_triggerbox_measurement_error,
996                sleep_dur,
997            )
998            .await
999            .map_err(|e| eyre::eyre!("on TriggerboxDevice::new: {e} {e:?}"))?;
1000            let query_dt2 = *query_dt;
1001            debug!("starting triggerbox task {}:{}", file!(), line!());
1002            let fut = async move {
1003                let result = triggerbox.run_forever(query_dt2).await;
1004                debug!("triggerbox task done {}:{}", file!(), line!());
1005                if let Err(e) = result {
1006                    error!("triggerbox result: {:?}", e);
1007                }
1008            };
1009            let _join_handle = tokio::spawn(fut);
1010        }
1011        TriggerType::FakeSync(FakeSyncConfig { framerate }) => {
1012            info!("No triggerbox configuration. Using fake synchronization.");
1013
1014            signal_triggerbox_connected.store(true, Ordering::SeqCst);
1015
1016            let mut expected_framerate = expected_framerate_arc.write().unwrap();
1017            *expected_framerate = Some(*framerate as f32);
1018
1019            let gain = 1.0 / framerate;
1020
1021            let now: chrono::DateTime<chrono::Utc> = chrono::Utc::now();
1022            let offset = strand_datetime_conversion::datetime_to_f64(&now);
1023
1024            (on_new_clock_model)(Some(braid_triggerbox::ClockModel {
1025                gain,
1026                n_measurements: 0,
1027                offset,
1028                residuals: 0.0,
1029            }));
1030        }
1031        TriggerType::PtpSync(ptpcfg) => {
1032            signal_triggerbox_connected.store(true, Ordering::SeqCst);
1033
1034            if let Some(periodic_signal_period_usec) = ptpcfg.periodic_signal_period_usec {
1035                let framerate = 1e6 / periodic_signal_period_usec;
1036                let mut expected_framerate = expected_framerate_arc.write().unwrap();
1037                *expected_framerate = Some(framerate as f32);
1038            }
1039        }
1040        TriggerType::DeviceTimestamp => {
1041            signal_triggerbox_connected.store(true, Ordering::SeqCst);
1042        }
1043    };
1044
1045    let expected_framerate_arc9 = expected_framerate_arc.clone();
1046
1047    let live_stats_collector = LiveStatsCollector::new(tracker.clone());
1048    let tracker2 = tracker.clone();
1049
1050    // decode UDP frames
1051    let raw_cam_data_stream =
1052        tokio_util::udp::UdpFramed::new(camdata_socket, CborPacketCodec::default());
1053
1054    // Initiate camera synchronization on startup
1055    let sync_pulse_pause_started_arc2 = sync_pulse_pause_started_arc.clone();
1056    let time_model_arc2 = time_model_arc.clone();
1057    let cam_manager2 = cam_manager.clone();
1058    let valve2 = valve.clone();
1059    let triggerbox_cmd2 = triggerbox_cmd.clone();
1060    let fake_sync = matches!(trigger_cfg, TriggerType::FakeSync(_));
1061    let _sync_start_jh = tokio::spawn(async move {
1062        let interval_stream = tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(
1063            std::time::Duration::from_secs(1),
1064        ));
1065        let mut interval_stream = valve2.wrap(interval_stream);
1066
1067        while let Some(_now) = interval_stream.next().await {
1068            let have_triggerbox = signal_triggerbox_connected.load(Ordering::SeqCst);
1069            let have_all_cameras = signal_all_cams_present.load(Ordering::SeqCst);
1070
1071            if have_triggerbox && have_all_cameras {
1072                info!("have triggerbox and all cameras. Synchronizing cameras.");
1073                synchronize_cameras(
1074                    triggerbox_cmd2,
1075                    fake_sync,
1076                    sync_pulse_pause_started_arc2.clone(),
1077                    cam_manager2.clone(),
1078                    time_model_arc2.clone(),
1079                )
1080                .await
1081                .unwrap();
1082                break;
1083            }
1084        }
1085    });
1086
1087    // Signal cameras are synchronized
1088
1089    let valve2 = valve.clone();
1090    let _sync_done_jh = tokio::spawn(async move {
1091        let interval_stream = tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(
1092            std::time::Duration::from_secs(1),
1093        ));
1094        let mut interval_stream = valve2.wrap(interval_stream);
1095        while let Some(_now) = interval_stream.next().await {
1096            let sync_done = signal_all_cams_synced.load(Ordering::SeqCst);
1097            if sync_done {
1098                info!("All cameras done synchronizing.");
1099
1100                // Send message to listeners.
1101                let mut tracker = shared_store.write().unwrap();
1102                tracker.modify(|shared| shared.all_expected_cameras_are_synced = true);
1103                break;
1104            }
1105        }
1106    });
1107
1108    let strand_cam_http_session_handler2 = strand_cam_http_session_handler.clone();
1109    let cam_manager2 = cam_manager.clone();
1110    let live_stats_collector2 = live_stats_collector.clone();
1111
1112    let packet_filter = move |r| {
1113        let live_stats_collector2 = live_stats_collector2.clone();
1114        let trigger_cfg = trigger_cfg.clone();
1115        let strand_cam_http_session_handler2 = strand_cam_http_session_handler2.clone();
1116        let cam_manager2 = cam_manager2.clone();
1117        let sync_pulse_pause_started_arc = sync_pulse_pause_started_arc.clone();
1118        let cam_manager = cam_manager.clone();
1119        let time_model_arc = time_model_arc.clone();
1120        async move {
1121            // vvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvvv
1122            // Start of closure for on each incoming packet.
1123
1124            // We run this closure for each incoming packet.
1125
1126            // Let's be sure about the type of our input.
1127            let r: std::result::Result<
1128                (braid_types::FlydraRawUdpPacket, std::net::SocketAddr),
1129                std::io::Error,
1130            > = r;
1131
1132            let (packet, _addr) = match r {
1133                Ok(r) => r,
1134                Err(e) => {
1135                    error!("{}", e);
1136                    return Some(StreamItem::EOF);
1137                }
1138            };
1139
1140            let raw_cam_name = RawCamName::new(packet.cam_name.clone());
1141            live_stats_collector2.register_new_frame_data(&raw_cam_name, packet.points.len());
1142
1143            // Create closure which is called only if there is a new frame offset
1144            // (which occurs upon synchronization).
1145            let send_new_frame_offset = |frame| {
1146                let strand_cam_http_session_handler = strand_cam_http_session_handler2.clone();
1147                let cam_name = raw_cam_name.clone();
1148                let fut_no_err = async move {
1149                    match strand_cam_http_session_handler
1150                        .send_frame_offset(&cam_name, frame)
1151                        .await
1152                    {
1153                        Ok(_) => {}
1154                        Err(e) => {
1155                            error!("Error sending frame offset: {}", e);
1156                        }
1157                    };
1158                };
1159                tokio::spawn(fut_no_err);
1160            };
1161
1162            let synced_frame = cam_manager2.got_new_frame_live(
1163                &packet,
1164                &sync_pulse_pause_started_arc,
1165                send_new_frame_offset,
1166                &trigger_cfg,
1167            );
1168
1169            let cam_num = cam_manager.cam_num(&raw_cam_name);
1170
1171            let cam_num = match cam_num {
1172                Some(cam_num) => cam_num,
1173                None => {
1174                    let known_raw_cam_names = cam_manager.all_raw_cam_names();
1175                    let cam_names = known_raw_cam_names
1176                        .iter()
1177                        .map(|x| format!("\"{}\"", x.as_str()))
1178                        .collect::<Vec<_>>()
1179                        .join(", ");
1180                    debug!(
1181                        "Unknown camera name \"{}\" ({} expected cameras: [{}]).",
1182                        raw_cam_name.as_str(),
1183                        known_raw_cam_names.len(),
1184                        cam_names
1185                    );
1186                    // Cannot compute cam_num, drop this data.
1187                    return None;
1188                }
1189            };
1190
1191            let (synced_frame, trigger_timestamp) = match synced_frame {
1192                Some(synced_frame) => {
1193                    let trigger_timestamp = match &trigger_cfg {
1194                        TriggerType::TriggerboxV1(_) => {
1195                            let time_model = time_model_arc.read().unwrap();
1196                            compute_trigger_timestamp(&time_model, synced_frame)
1197                        }
1198                        TriggerType::FakeSync(_) => {
1199                            // There is no trigger clock. The camera host clock
1200                            // is the best available approximation of the
1201                            // acquisition time (and, with fake sync, plays the
1202                            // role the triggerbox clock model otherwise would).
1203                            Some(FlydraFloatTimestampLocal::from_f64(
1204                                packet.cam_received_time.as_f64(),
1205                            ))
1206                        }
1207                        TriggerType::PtpSync(_) => {
1208                            // In case where we trust camera sync data, use
1209                            // timestamp from camera. All packets from all
1210                            // cameras should have this same timestamp, so it
1211                            // shouldn't matter which camera we use.
1212                            packet.device_timestamp.map(|device_timestamp| {
1213                                let ptp_stamp = braid_types::PtpStamp::new(device_timestamp);
1214                                let device_timestamp_chrono =
1215                                    chrono::DateTime::<chrono::Utc>::try_from(ptp_stamp.clone())
1216                                        .unwrap();
1217                                device_timestamp_chrono.into()
1218                            })
1219                        }
1220                        TriggerType::DeviceTimestamp => {
1221                            todo!();
1222                        }
1223                    };
1224                    (synced_frame, trigger_timestamp)
1225                }
1226                None => {
1227                    // cannot compute synced_frame number, drop this data
1228                    return None;
1229                }
1230            };
1231
1232            let frame_data = flydra2::FrameData::new(
1233                raw_cam_name,
1234                cam_num,
1235                synced_frame,
1236                trigger_timestamp,
1237                packet.cam_received_time,
1238                packet.device_timestamp,
1239                packet.block_id,
1240            );
1241
1242            assert!(packet.points.len() < u8::MAX as usize);
1243            let points = packet
1244                .points
1245                .into_iter()
1246                .enumerate()
1247                .map(|(idx, pt)| {
1248                    assert!(idx <= 255);
1249                    flydra2::NumberedRawUdpPoint { idx: idx as u8, pt }
1250                })
1251                .collect();
1252
1253            let fdp = FrameDataAndPoints { frame_data, points };
1254            Some(StreamItem::Packet(fdp))
1255            // This is the end of closure for each incoming packet.
1256            // ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
1257        }
1258    };
1259
1260    let flydra2_stream = raw_cam_data_stream.filter_map(packet_filter);
1261
1262    let (data_tx, data_rx) = tokio::sync::mpsc::channel(50);
1263
1264    let model_pose_server_addr = mainbrain_config.model_server_addr;
1265    tokio::spawn(flydra2::new_model_server(data_rx, model_pose_server_addr));
1266
1267    {
1268        let mut tracker = tracker2.write().unwrap();
1269        tracker.modify(|shared| shared.model_server_addr = Some(model_pose_server_addr))
1270    }
1271
1272    let expected_framerate: Option<f32> = *expected_framerate_arc9.read().unwrap();
1273    info!("expected_framerate: {:?}", expected_framerate);
1274
1275    coord_processor.add_listener(data_tx);
1276    let coord_proc_fut = coord_processor.consume_stream(flydra2_stream, expected_framerate);
1277
1278    // We "block" (in an async way) here for the entire runtime of the program.
1279    // The first one of these to exit will end all of them. This should be
1280    // `coord_proc_fut`.
1281    tokio::select! {
1282        _ = event_broadcast_fut => {
1283            info!("Event broadcaster finished.");
1284        },
1285        http_serve_res = http_serve_future => {
1286            info!("HTTP Server finished.");
1287            http_serve_res?;
1288        },
1289        strand_cam_set_opt_res = strand_cam_set.join_next() => {
1290            // One of the strand camera tasks has finished.
1291            info!("Strand Camera future set member finished.");
1292            if let Some(strand_cam_res) = strand_cam_set_opt_res {
1293                strand_cam_res?;
1294            }
1295        },
1296        res_writer_jh = coord_proc_fut => {
1297            info!("Coordinate processor finished.");
1298            // Allow writer task time to finish writing.
1299            debug!("Runtime ending. Joining coord_processor.consume_stream future.");
1300            res_writer_jh?.await??;
1301        },
1302    };
1303    // We reach here after the first of the above futures finishes. We now drop
1304    // the other futures, thus cancelling them.
1305
1306    debug!("braid-run finishing.");
1307
1308    Ok(())
1309}
1310
1311#[derive(Clone)]
1312struct LiveStatsCollector {
1313    shared: SharedStore,
1314    collected: Arc<RwLock<BTreeMap<RawCamName, LiveStatsAccum>>>,
1315}
1316
1317#[derive(Debug)]
1318struct LiveStatsAccum {
1319    start: std::time::Instant,
1320    n_frames: usize,
1321    n_points: usize,
1322}
1323
1324impl LiveStatsAccum {
1325    fn new() -> Self {
1326        Self {
1327            start: std::time::Instant::now(),
1328            n_frames: 0,
1329            n_points: 0,
1330        }
1331    }
1332    fn update(&mut self, n_points: usize) {
1333        self.n_frames += 1;
1334        self.n_points += n_points;
1335    }
1336    fn get_results_and_reset(&mut self) -> braid_types::RecentStats {
1337        let elapsed = self.start.elapsed().as_secs_f64();
1338        let fps = if elapsed > 0.0 {
1339            self.n_frames as f64 / elapsed
1340        } else {
1341            0.0
1342        };
1343        let recent = braid_types::RecentStats {
1344            total_frames_collected: 0,
1345            frames_collected: self.n_frames,
1346            fps,
1347            points_detected: self.n_points,
1348        };
1349        self.start = std::time::Instant::now();
1350        self.n_frames = 0;
1351        self.n_points = 0;
1352        recent
1353    }
1354}
1355
1356impl LiveStatsCollector {
1357    fn new(shared: SharedStore) -> Self {
1358        let collected = Arc::new(RwLock::new(BTreeMap::new()));
1359        Self { shared, collected }
1360    }
1361
1362    fn register_new_frame_data(&self, name: &RawCamName, n_points: usize) {
1363        let to_send = {
1364            // scope for lock on self.collected
1365            let mut collected = self.collected.write().unwrap();
1366            let entry = collected
1367                .entry(name.clone())
1368                .or_insert_with(LiveStatsAccum::new);
1369            entry.update(n_points);
1370
1371            if entry.start.elapsed() > std::time::Duration::from_secs(1) {
1372                Some((name.clone(), entry.get_results_and_reset()))
1373            } else {
1374                None
1375            }
1376        };
1377        if let Some((name, recent_stats)) = to_send {
1378            // scope for shared scope
1379            let mut tracker = self.shared.write().unwrap();
1380            tracker.modify(|shared| {
1381                for cc in shared.connected_cameras.iter_mut() {
1382                    if cc.name == name {
1383                        let old_total = cc.recent_stats.total_frames_collected;
1384                        cc.recent_stats = recent_stats.clone();
1385                        cc.recent_stats.total_frames_collected =
1386                            old_total + recent_stats.frames_collected;
1387                        break;
1388                    }
1389                }
1390            });
1391        }
1392    }
1393}
1394
1395pub(crate) async fn toggle_saving_csv_tables(
1396    start_saving: bool,
1397    expected_framerate_arc: Arc<RwLock<Option<f32>>>,
1398    output_base_dirname: std::path::PathBuf,
1399    braidz_write_tx_weak: tokio::sync::mpsc::WeakSender<flydra2::SaveToDiskMsg>,
1400    per_cam_data_arc: Arc<RwLock<BTreeMap<RawCamName, PerCamSaveData>>>,
1401    shared_data: SharedStore,
1402) {
1403    if start_saving {
1404        let expected_framerate: Option<f32> = *expected_framerate_arc.read().unwrap();
1405        let local: chrono::DateTime<chrono::Local> = chrono::Local::now();
1406        let dirname = local.format("%Y%m%d_%H%M%S.braid").to_string();
1407        let mut my_dir = output_base_dirname.clone();
1408        my_dir.push(dirname);
1409        let per_cam_data = {
1410            // small scope for read lock
1411            let per_cam_data_ref = per_cam_data_arc.read().unwrap();
1412            (*per_cam_data_ref).clone()
1413        };
1414        let cfg = flydra2::StartSavingCsvConfig {
1415            out_dir: my_dir.clone(),
1416            local: Some(local),
1417            git_rev: env!("GIT_HASH").to_string(),
1418            fps: expected_framerate,
1419            per_cam_data,
1420            print_stats: false,
1421            save_performance_histograms: true,
1422        };
1423
1424        match braidz_write_tx_weak.upgrade() {
1425            Some(braidz_write_tx) => {
1426                // `braidz_write_tx` will be dropped after this scope.
1427                braidz_write_tx
1428                    .send(flydra2::SaveToDiskMsg::StartSavingCsv(cfg))
1429                    .await
1430                    .unwrap();
1431                info!("saving data to \"{}\"", my_dir.display());
1432            }
1433            _ => {
1434                error!("data writing thread lost. Not saving data as requested");
1435            }
1436        }
1437
1438        {
1439            let mut tracker = shared_data.write().unwrap();
1440            tracker.modify(|store| {
1441                store.csv_tables_dirname = Some(RecordingPath::new(my_dir.display().to_string()));
1442            });
1443        }
1444    } else {
1445        match braidz_write_tx_weak.upgrade() {
1446            Some(braidz_write_tx) => {
1447                // `braidz_write_tx` will be dropped after this scope.
1448                braidz_write_tx
1449                    .send(flydra2::SaveToDiskMsg::StopSavingCsv)
1450                    .await
1451                    .unwrap_or(()); // ignore error on shutdown
1452                info!("stopping saving");
1453            }
1454            _ => {
1455                error!("data writing thread lost. Could not stop saving data as requested");
1456            }
1457        }
1458
1459        {
1460            let mut tracker = shared_data.write().unwrap();
1461            tracker.modify(|store| {
1462                store.csv_tables_dirname = None;
1463            });
1464        }
1465    }
1466}
1467
1468async fn synchronize_cameras(
1469    triggerbox_cmd: Option<tokio::sync::mpsc::Sender<braid_triggerbox::Cmd>>,
1470    fake_sync: bool,
1471    sync_pulse_pause_started_arc: Arc<RwLock<Option<flydra2::SyncStart>>>,
1472    mut cam_manager: flydra2::ConnectedCamerasManager,
1473    time_model_arc: Arc<RwLock<Option<strand_cam_bui_types::ClockModel>>>,
1474) -> Result<()> {
1475    info!("preparing to synchronize cameras");
1476
1477    // This time must be prior to actually resetting sync data.
1478    {
1479        let mut sync_pulse_pause_started = sync_pulse_pause_started_arc.write().unwrap();
1480        *sync_pulse_pause_started = Some(flydra2::SyncStart::now());
1481    }
1482
1483    // Now we can reset the sync data.
1484    cam_manager.reset_sync_data();
1485
1486    {
1487        let mut guard = time_model_arc.write().unwrap();
1488        *guard = None;
1489    }
1490
1491    if let Some(tx) = triggerbox_cmd {
1492        begin_cam_sync_triggerbox_in_process(tx).await?;
1493    }
1494
1495    if fake_sync {
1496        info!("Using fake synchronization method.");
1497    }
1498    Ok(())
1499}
1500
1501async fn begin_cam_sync_triggerbox_in_process(
1502    tx: tokio::sync::mpsc::Sender<braid_triggerbox::Cmd>,
1503) -> Result<()> {
1504    // This is the case when the triggerbox is within this process.
1505    info!("preparing for triggerbox to temporarily stop sending pulses");
1506
1507    info!("requesting triggerbox to stop sending pulses");
1508    use braid_triggerbox::Cmd::*;
1509    tx.send(StopPulsesAndReset).await?;
1510    tokio::time::sleep(std::time::Duration::from_secs(TRIGGERBOX_SYNC_SECONDS)).await;
1511    tx.send(StartPulses).await?;
1512    info!("requesting triggerbox to start sending pulses again");
1513    Ok(())
1514}
1515
1516fn to_event_frame(state: &BraidHttpApiSharedState) -> String {
1517    let buf = serde_json::to_string(&state).unwrap();
1518    let frame_string = format!("event: {BRAID_EVENT_NAME}\ndata: {buf}\n\n");
1519    frame_string
1520}
1521
1522/// Build the Server-Sent Events frame announcing that the server is quitting.
1523///
1524/// The data payload is unused by the frontend (the event name alone is the
1525/// signal), but SSE frames must carry a `data:` line.
1526fn quit_event_frame() -> String {
1527    format!("event: {BRAID_QUIT_EVENT_NAME}\ndata: quit\n\n")
1528}