livekit_datatrack/local/
mod.rs1use crate::{
16 api::{DataTrack, DataTrackFrame, DataTrackInfo, DataTrackSchemaError, InternalError},
17 schema::{DataTrackFrameEncoding, DataTrackSchemaId},
18 track::DataTrackInner,
19};
20use std::{fmt, marker::PhantomData, sync::Arc};
21use thiserror::Error;
22use tokio::sync::{mpsc, watch};
23
24pub(crate) mod events;
25pub(crate) mod manager;
26pub(crate) mod proto;
27
28mod packetizer;
29mod pipeline;
30
31pub type LocalDataTrack = DataTrack<Local>;
33
34#[derive(Debug, Clone)]
39pub struct Local;
40
41impl DataTrack<Local> {
42 pub(crate) fn new(info: Arc<DataTrackInfo>, inner: LocalTrackInner) -> Self {
43 Self { info, inner: Arc::new(inner.into()), _location: PhantomData }
44 }
45
46 fn inner(&self) -> &LocalTrackInner {
47 match &*self.inner {
48 DataTrackInner::Local(track) => track,
49 DataTrackInner::Remote(_) => unreachable!(), }
51 }
52}
53
54impl DataTrack<Local> {
55 pub fn try_push(&self, frame: DataTrackFrame) -> Result<(), PushFrameError> {
85 match self.inner().publish_state() {
86 manager::PublishState::Republishing => {
87 return Err(PushFrameError::new(frame, PushFrameErrorReason::QueueFull))?
88 }
89 manager::PublishState::Unpublished => {
90 return Err(PushFrameError::new(frame, PushFrameErrorReason::TrackUnpublished))?;
91 }
92 manager::PublishState::Published => {}
93 }
94 self.inner()
95 .frame_tx
96 .try_send(frame)
97 .map_err(|err| PushFrameError::new(err.into_inner(), PushFrameErrorReason::QueueFull))
98 }
99
100 pub fn unpublish(&self) {
102 self.inner().local_unpublish();
103 }
104}
105
106#[derive(Debug, Clone)]
107pub(crate) struct LocalTrackInner {
108 pub frame_tx: mpsc::Sender<DataTrackFrame>,
109 pub state_tx: watch::Sender<manager::PublishState>,
110}
111
112impl LocalTrackInner {
113 fn publish_state(&self) -> manager::PublishState {
114 *self.state_tx.borrow()
115 }
116
117 pub(crate) fn is_published(&self) -> bool {
118 self.publish_state() != manager::PublishState::Unpublished
121 }
122
123 pub(crate) async fn wait_for_unpublish(&self) {
124 _ = self
125 .state_tx
126 .subscribe()
127 .wait_for(|state| *state == manager::PublishState::Unpublished)
128 .await
129 }
130
131 fn local_unpublish(&self) {
132 _ = self.state_tx.send(manager::PublishState::Unpublished);
133 }
134}
135
136impl Drop for LocalTrackInner {
137 fn drop(&mut self) {
138 self.local_unpublish();
140 }
141}
142
143#[derive(Clone, Debug)]
155pub struct DataTrackOptions {
156 pub(crate) name: String,
157 pub(crate) schema: Option<DataTrackSchemaId>,
158 pub(crate) frame_encoding: Option<DataTrackFrameEncoding>,
159}
160
161impl DataTrackOptions {
162 pub fn new(name: impl Into<String>) -> Self {
171 Self { name: name.into(), schema: None, frame_encoding: None }
172 }
173
174 pub fn with_schema(self, schema: DataTrackSchemaId) -> Self {
176 Self { schema: Some(schema), ..self }
177 }
178
179 pub fn with_frame_encoding(self, encoding: DataTrackFrameEncoding) -> Self {
181 Self { frame_encoding: Some(encoding), ..self }
182 }
183}
184
185impl From<String> for DataTrackOptions {
186 fn from(name: String) -> Self {
187 Self::new(name)
188 }
189}
190
191impl From<&str> for DataTrackOptions {
192 fn from(name: &str) -> Self {
193 Self::new(name.to_string())
194 }
195}
196
197#[derive(Debug, Error)]
199#[cfg_attr(feature = "uniffi", derive(uniffi::Error))]
200#[cfg_attr(feature = "uniffi", uniffi(flat_error))]
201pub enum PublishError {
202 #[error("Data track publishing unauthorized")]
207 NotAllowed,
208
209 #[error("Track name already taken")]
211 DuplicateName,
212
213 #[error("Track name invalid")]
218 InvalidName,
219
220 #[error("Publish data track timed-out")]
222 Timeout,
223
224 #[error("Data track publication limit reached")]
226 LimitReached,
227
228 #[error("Room disconnected")]
230 Disconnected,
231
232 #[error(transparent)]
234 InvalidSchema(DataTrackSchemaError),
235
236 #[error(transparent)]
238 Internal(#[from] InternalError),
239}
240
241#[derive(Debug, Error)]
243#[error("Failed to publish frame: {reason}")]
244pub struct PushFrameError {
245 frame: DataTrackFrame,
246 reason: PushFrameErrorReason,
247}
248
249impl PushFrameError {
250 pub(crate) fn new(frame: DataTrackFrame, reason: PushFrameErrorReason) -> Self {
251 Self { frame, reason }
252 }
253
254 pub fn reason(&self) -> PushFrameErrorReason {
256 self.reason
257 }
258
259 pub fn into_frame(self) -> DataTrackFrame {
264 self.frame
265 }
266}
267
268#[derive(Debug, Clone, Copy)]
270#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
271pub enum PushFrameErrorReason {
272 TrackUnpublished,
274 QueueFull,
276}
277
278impl fmt::Display for PushFrameErrorReason {
279 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
280 match self {
281 Self::TrackUnpublished => write!(f, "track unpublished"),
282 Self::QueueFull => write!(f, "queue full"),
283 }
284 }
285}