Skip to main content

o_sfu_core/engine/media_transport/
teardown.rs

1use std::ptr;
2
3use tracing::warn;
4
5use super::{
6    MediaTransport, TransportAdapterError, TransportMediaId, TransportRelayRouteAction,
7    TransportSessionKey, TransportSourceKey, rtc::RtcWorkerCommand,
8};
9use crate::engine::MediaWorkerId;
10
11#[derive(Debug)]
12pub(crate) enum TransportTeardown {
13    CloseSession {
14        session_key: TransportSessionKey,
15    },
16    RemoveMedia {
17        session_key: TransportSessionKey,
18        transport_media_id: TransportMediaId,
19    },
20    ReleaseRelayRoute {
21        source: TransportSourceKey,
22        target_media_worker_id: MediaWorkerId,
23    },
24}
25
26impl TransportTeardown {
27    pub(crate) fn session_key(&self) -> &TransportSessionKey {
28        match self {
29            Self::CloseSession { session_key } | Self::RemoveMedia { session_key, .. } => {
30                session_key
31            }
32            Self::ReleaseRelayRoute { source, .. } => source.session_key(),
33        }
34    }
35}
36
37impl MediaTransport {
38    /// Runs teardown idempotently and continues after terminal failures.
39    ///
40    /// Each failed media or relay item escalates once to its session because
41    /// partial worker state may remain. Errors are consumed so one failure cannot
42    /// suppress later cleanup.
43    pub(crate) async fn teardown(&self, teardowns: impl IntoIterator<Item = TransportTeardown>) {
44        for teardown in teardowns {
45            let (session_key, result, transport_media_id, target_media_worker_id) = match &teardown
46            {
47                TransportTeardown::CloseSession { session_key } => (
48                    session_key,
49                    self.close_session(session_key).await,
50                    None,
51                    None,
52                ),
53                TransportTeardown::RemoveMedia {
54                    session_key,
55                    transport_media_id,
56                } => (
57                    session_key,
58                    self.remove_media(session_key, *transport_media_id).await,
59                    Some(*transport_media_id),
60                    None,
61                ),
62                TransportTeardown::ReleaseRelayRoute {
63                    source,
64                    target_media_worker_id,
65                } => (
66                    source.session_key(),
67                    self.release_relay_route(source, *target_media_worker_id)
68                        .await,
69                    Some(source.transport_media_id()),
70                    Some(*target_media_worker_id),
71                ),
72            };
73            let Err(error) = result else {
74                continue;
75            };
76            self.metrics.record_transport_cleanup_failure();
77            warn!(
78                ?session_key,
79                ?transport_media_id,
80                ?target_media_worker_id,
81                ?error,
82                "media transport teardown reached terminal failure"
83            );
84            // `RemoveMedia` or `ReleaseRelayRoute` failure may leave worker state behind.
85            // Escalate once to `CloseSession`, whose own failure is already terminal.
86            if transport_media_id.is_some()
87                && let Err(fallback_error) = self.close_session(session_key).await
88            {
89                warn!(
90                    ?session_key,
91                    ?transport_media_id,
92                    ?target_media_worker_id,
93                    ?fallback_error,
94                    "media transport teardown session fallback failed"
95                );
96            }
97        }
98    }
99
100    async fn close_session(
101        &self,
102        session_key: &TransportSessionKey,
103    ) -> Result<(), TransportAdapterError> {
104        let worker = self.require_worker_for_user(session_key)?;
105        worker
106            .request_worker(|response| RtcWorkerCommand::CloseSession {
107                session_key: session_key.clone(),
108                response,
109            })
110            .await
111    }
112
113    pub(super) async fn remove_media(
114        &self,
115        session_key: &TransportSessionKey,
116        transport_media_id: TransportMediaId,
117    ) -> Result<(), TransportAdapterError> {
118        let worker = self.require_worker_for_user(session_key)?;
119        worker
120            .request_worker(|response| RtcWorkerCommand::RemoveMedia {
121                session_key: session_key.clone(),
122                transport_media_id,
123                response,
124            })
125            .await
126    }
127
128    async fn release_relay_route(
129        &self,
130        source: &TransportSourceKey,
131        target_media_worker_id: MediaWorkerId,
132    ) -> Result<(), TransportAdapterError> {
133        let source_worker = self.require_worker_for_user(source.session_key())?;
134        let target_worker = self.require_worker_for_media_worker_id(target_media_worker_id)?;
135        if ptr::eq(source_worker, target_worker) {
136            return Ok(());
137        }
138        let request =
139            target_worker.relay_route_request(source.clone(), TransportRelayRouteAction::Release);
140        source_worker
141            .request_worker(|response| RtcWorkerCommand::RouteControl {
142                request,
143                response: Some(response),
144            })
145            .await
146    }
147}