1use std::time::Duration;
2
3use serde::{Deserialize, Serialize};
4use thiserror::Error;
5use url::Url;
6
7use crate::{DecoderBitstreamFormat, NativeFrameColorMetadata, NativeFrameHdrMetadata};
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default, Serialize, Deserialize)]
11pub enum SourceNormalizerNormalizeLevel {
12 #[default]
14 #[serde(alias = "remux_only", alias = "remux-only")]
15 RemuxOnly = 1,
16 #[serde(alias = "packet_repair", alias = "packet-repair")]
18 PacketRepair = 2,
19}
20
21#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
23pub struct SourceNormalizerRequiredCapabilities {
24 pub libraries: Vec<String>,
25 pub demuxers: Vec<String>,
26 pub muxers: Vec<String>,
27 pub protocols: Vec<String>,
28 pub parsers: Vec<String>,
29 pub bitstream_filters: Vec<String>,
30 #[serde(default)]
31 pub tls: Option<String>,
32 #[serde(default)]
33 pub network: bool,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
38pub struct SourceNormalizerPacketCapabilities {
39 pub supported_runtime_profiles: Vec<String>,
40 pub max_level: SourceNormalizerNormalizeLevel,
41 pub media_kinds: Vec<SourceNormalizerPacketMediaKind>,
42 pub codecs: Vec<String>,
43 pub bitstream_formats: Vec<DecoderBitstreamFormat>,
44 pub supports_seek: bool,
45 pub supports_flush: bool,
46 pub required_capabilities: SourceNormalizerRequiredCapabilities,
47 pub max_sessions: Option<u32>,
48}
49
50impl SourceNormalizerPacketCapabilities {
51 pub fn supports_runtime_profile(&self, runtime_profile: &str) -> bool {
53 self.supported_runtime_profiles
54 .iter()
55 .any(|profile| profile.eq_ignore_ascii_case(runtime_profile))
56 }
57
58 pub fn supports_codec(&self, codec: &str) -> bool {
60 self.codecs
61 .iter()
62 .any(|candidate| candidate.eq_ignore_ascii_case(codec))
63 }
64
65 pub fn supports_media_kind(&self, media_kind: SourceNormalizerPacketMediaKind) -> bool {
67 self.media_kinds.contains(&media_kind)
68 }
69
70 pub fn supports_bitstream_format(&self, format: &DecoderBitstreamFormat) -> bool {
72 self.bitstream_formats.contains(format)
73 }
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
78#[serde(rename_all = "camelCase")]
79pub enum SourceNormalizerOutputRoute {
80 Fmp4LocalStream,
82 HlsShortWindow,
84 PacketStream,
86}
87
88impl SourceNormalizerOutputRoute {
89 pub fn wire_name(self) -> &'static str {
90 match self {
91 Self::Fmp4LocalStream => "fmp4LocalStream",
92 Self::HlsShortWindow => "hlsShortWindow",
93 Self::PacketStream => "packetStream",
94 }
95 }
96}
97
98#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
100pub struct SourceNormalizerResourceCachePolicy {
101 pub session_read_buffer_bytes: u64,
103 pub manifest_snapshot_bytes: u64,
105 pub session_disk_soft_cap_bytes: u64,
107 pub global_disk_soft_cap_bytes: u64,
109}
110
111impl Default for SourceNormalizerResourceCachePolicy {
112 fn default() -> Self {
113 Self {
114 session_read_buffer_bytes: 4 * 1024 * 1024,
115 manifest_snapshot_bytes: 512 * 1024,
116 session_disk_soft_cap_bytes: 512 * 1024 * 1024,
117 global_disk_soft_cap_bytes: 1536 * 1024 * 1024,
118 }
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
124pub struct SourceNormalizerResourceCapabilities {
125 pub supported_runtime_profiles: Vec<String>,
126 pub supported_output_routes: Vec<SourceNormalizerOutputRoute>,
127 pub max_level: SourceNormalizerNormalizeLevel,
128 pub content_types: Vec<String>,
129 pub supports_growing_resources: bool,
130 pub supports_range_reads: bool,
131 pub supports_cancel: bool,
132 pub required_capabilities: SourceNormalizerRequiredCapabilities,
133 pub cache_policy: SourceNormalizerResourceCachePolicy,
134 pub max_sessions: Option<u32>,
135}
136
137impl SourceNormalizerResourceCapabilities {
138 pub fn supports_runtime_profile(&self, runtime_profile: &str) -> bool {
140 self.supported_runtime_profiles
141 .iter()
142 .any(|profile| profile.eq_ignore_ascii_case(runtime_profile))
143 }
144
145 pub fn supports_output_route(&self, route: SourceNormalizerOutputRoute) -> bool {
147 self.supported_output_routes.contains(&route)
148 }
149}
150
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
153pub enum SourceNormalizerSessionRequirements {
154 Packet(SourceNormalizerPacketSessionRequirements),
155 Resource(SourceNormalizerResourceSessionRequirements),
156}
157
158impl SourceNormalizerSessionRequirements {
159 pub fn missing_capabilities(
161 &self,
162 capabilities: &SourceNormalizerSessionCapabilities<'_>,
163 ) -> Vec<String> {
164 match (self, capabilities) {
165 (
166 Self::Packet(requirements),
167 SourceNormalizerSessionCapabilities::Packet(capabilities),
168 ) => requirements.missing_capabilities(capabilities),
169 (
170 Self::Resource(requirements),
171 SourceNormalizerSessionCapabilities::Resource(capabilities),
172 ) => requirements.missing_capabilities(capabilities),
173 (Self::Packet(_), SourceNormalizerSessionCapabilities::Resource(_)) => {
174 vec!["packet stream route".to_owned()]
175 }
176 (Self::Resource(_), SourceNormalizerSessionCapabilities::Packet(_)) => {
177 vec!["resource output route".to_owned()]
178 }
179 }
180 }
181}
182
183#[derive(Debug, Clone, Copy)]
185pub enum SourceNormalizerSessionCapabilities<'a> {
186 Packet(&'a SourceNormalizerPacketCapabilities),
187 Resource(&'a SourceNormalizerResourceCapabilities),
188}
189
190#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192pub struct SourceNormalizerPacketSessionRequirements {
193 pub runtime_profile: String,
194 #[serde(default)]
195 pub media_kind: Option<SourceNormalizerPacketMediaKind>,
196 #[serde(default)]
197 pub codec: Option<String>,
198 #[serde(default)]
199 pub bitstream_format: Option<DecoderBitstreamFormat>,
200 #[serde(default)]
201 pub require_seek: bool,
202 #[serde(default)]
203 pub require_flush: bool,
204 #[serde(default)]
205 pub require_lease_cleanup: bool,
206}
207
208impl SourceNormalizerPacketSessionRequirements {
209 pub fn native_video(runtime_profile: impl Into<String>, codec: impl Into<String>) -> Self {
211 Self {
212 runtime_profile: runtime_profile.into(),
213 media_kind: Some(SourceNormalizerPacketMediaKind::Video),
214 codec: Some(codec.into()),
215 bitstream_format: None,
216 require_seek: false,
217 require_flush: true,
218 require_lease_cleanup: true,
219 }
220 }
221
222 pub fn missing_capabilities(
224 &self,
225 capabilities: &SourceNormalizerPacketCapabilities,
226 ) -> Vec<String> {
227 let mut missing = Vec::new();
228 if !self.runtime_profile.is_empty()
229 && !capabilities.supports_runtime_profile(&self.runtime_profile)
230 {
231 missing.push(format!("runtime profile {}", self.runtime_profile));
232 }
233 if let Some(media_kind) = self.media_kind
234 && !capabilities.supports_media_kind(media_kind)
235 {
236 missing.push(format!("packet media kind {media_kind:?}"));
237 }
238 if let Some(codec) = &self.codec
239 && !capabilities.supports_codec(codec)
240 {
241 missing.push(format!("packet codec {codec}"));
242 }
243 if let Some(format) = &self.bitstream_format
244 && !capabilities.supports_bitstream_format(format)
245 {
246 missing.push(format!("packet bitstream format {format:?}"));
247 }
248 if self.require_seek && !capabilities.supports_seek {
249 missing.push("packet seek support".to_owned());
250 }
251 if self.require_flush && !capabilities.supports_flush {
252 missing.push("packet flush support".to_owned());
253 }
254 if self.require_lease_cleanup && !capabilities.supports_flush {
255 missing.push("outstanding lease cleanup".to_owned());
256 }
257 missing
258 }
259}
260
261#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
263pub struct SourceNormalizerResourceSessionRequirements {
264 pub runtime_profile: String,
265 pub output_route: SourceNormalizerOutputRoute,
266 #[serde(default)]
267 pub content_type: Option<String>,
268 #[serde(default)]
269 pub require_growing_resources: bool,
270 #[serde(default)]
271 pub require_range_reads: bool,
272 #[serde(default)]
273 pub require_cancel: bool,
274}
275
276impl SourceNormalizerResourceSessionRequirements {
277 pub fn missing_capabilities(
279 &self,
280 capabilities: &SourceNormalizerResourceCapabilities,
281 ) -> Vec<String> {
282 let mut missing = Vec::new();
283 if !self.runtime_profile.is_empty()
284 && !capabilities.supports_runtime_profile(&self.runtime_profile)
285 {
286 missing.push(format!("runtime profile {}", self.runtime_profile));
287 }
288 if !capabilities.supports_output_route(self.output_route) {
289 missing.push(format!("resource output route {:?}", self.output_route));
290 }
291 if let Some(content_type) = &self.content_type {
292 let supported = capabilities
293 .content_types
294 .iter()
295 .any(|candidate| candidate.eq_ignore_ascii_case(content_type));
296 if !supported {
297 missing.push(format!("content type {content_type}"));
298 }
299 }
300 if self.require_growing_resources && !capabilities.supports_growing_resources {
301 missing.push("growing resources".to_owned());
302 }
303 if self.require_range_reads && !capabilities.supports_range_reads {
304 missing.push("range reads".to_owned());
305 }
306 if self.require_cancel && !capabilities.supports_cancel {
307 missing.push("cancel support".to_owned());
308 }
309 missing
310 }
311}
312
313#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default, Serialize, Deserialize)]
315pub enum SourceNormalizerPacketMediaKind {
316 #[default]
317 Video,
318 Audio,
319 Subtitle,
320}
321
322#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
324pub struct SourceNormalizerPacketSessionConfig {
325 pub runtime_profile: String,
326 pub input: String,
327 #[serde(default)]
328 pub headers: Vec<(String, String)>,
329 #[serde(default)]
330 pub startup_timeout_ms: Option<u64>,
331 #[serde(default)]
332 pub session_timeout_ms: Option<u64>,
333 #[serde(default)]
334 pub preferred_media_kind: SourceNormalizerPacketMediaKind,
335}
336
337#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
339pub struct SourceNormalizerResourceSessionConfig {
340 pub runtime_profile: String,
341 pub input: String,
342 #[serde(default)]
343 pub headers: Vec<(String, String)>,
344 pub output_root: String,
345 #[serde(default)]
346 pub cache_policy: SourceNormalizerResourceCachePolicy,
347 #[serde(default)]
348 pub preferred_route: Option<SourceNormalizerOutputRoute>,
349 #[serde(default)]
350 pub startup_timeout_ms: Option<u64>,
351 #[serde(default)]
352 pub read_idle_timeout_ms: Option<u64>,
353}
354
355#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
357pub struct SourceNormalizerPacketTrackInfo {
358 pub stream_index: u32,
359 pub media_kind: SourceNormalizerPacketMediaKind,
360 pub codec: String,
361 #[serde(default)]
362 pub extradata: Vec<u8>,
363 #[serde(default)]
364 pub bitstream_format: Option<DecoderBitstreamFormat>,
365 #[serde(default)]
366 pub width: Option<u32>,
367 #[serde(default)]
368 pub height: Option<u32>,
369 #[serde(default)]
370 pub coded_width: Option<u32>,
371 #[serde(default)]
372 pub coded_height: Option<u32>,
373 #[serde(default)]
374 pub reorder_depth: Option<u32>,
375 #[serde(default)]
376 pub sample_rate: Option<u32>,
377 #[serde(default)]
378 pub channels: Option<u16>,
379 #[serde(default)]
380 pub channel_layout: Option<String>,
381 #[serde(default)]
382 pub codec_delay_samples: Option<u32>,
383 #[serde(default)]
384 pub priming_samples: Option<u32>,
385 #[serde(default)]
386 pub trailing_padding_samples: Option<u32>,
387 #[serde(default)]
388 pub seek_preroll_samples: Option<u32>,
389 #[serde(default)]
390 pub color: Option<NativeFrameColorMetadata>,
391 #[serde(default)]
392 pub hdr: Option<NativeFrameHdrMetadata>,
393 #[serde(default)]
394 pub frame_rate: Option<f64>,
395 #[serde(default)]
396 pub time_base_num: Option<i32>,
397 #[serde(default)]
398 pub time_base_den: Option<i32>,
399}
400
401#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
403pub struct SourceNormalizerPacketStreamInfo {
404 #[serde(default)]
405 pub session_id: Option<String>,
406 #[serde(default)]
407 pub normalizer_name: Option<String>,
408 #[serde(default)]
409 pub runtime_profile: Option<String>,
410 #[serde(default)]
411 pub selected_backend: Option<String>,
412 pub tracks: Vec<SourceNormalizerPacketTrackInfo>,
413 #[serde(default)]
414 pub selected_track_index: Option<u32>,
415 #[serde(default)]
416 pub duration_millis: Option<u64>,
417 #[serde(default)]
418 pub seekable: bool,
419}
420
421#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
423pub struct SourceNormalizerResourceInfo {
424 pub role: String,
425 pub path: String,
426 #[serde(default)]
427 pub content_type: Option<String>,
428 #[serde(default)]
429 pub byte_length: Option<u64>,
430 #[serde(default)]
431 pub growing: bool,
432}
433
434#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
436pub struct SourceNormalizerResourceSessionInfo {
437 #[serde(default)]
438 pub session_id: Option<String>,
439 #[serde(default)]
440 pub normalizer_name: Option<String>,
441 #[serde(default)]
442 pub runtime_profile: Option<String>,
443 #[serde(default)]
444 pub selected_backend: Option<String>,
445 pub output_route: SourceNormalizerOutputRoute,
446 pub container: String,
447 #[serde(default)]
448 pub primary_resource_path: Option<String>,
449 #[serde(default)]
450 pub primary_content_type: Option<String>,
451 #[serde(default)]
452 pub resources: Vec<SourceNormalizerResourceInfo>,
453 #[serde(default)]
454 pub tracks: Vec<SourceNormalizerPacketTrackInfo>,
455 #[serde(default)]
456 pub duration_millis: Option<u64>,
457 #[serde(default)]
458 pub seekable: bool,
459 #[serde(default)]
460 pub disk_bytes_used: Option<u64>,
461}
462
463#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
465#[serde(rename_all = "camelCase")]
466pub enum SourceNormalizerResourceSessionState {
467 Starting,
468 Ready,
469 Running,
470 Completed,
471 Failed,
472 Cancelled,
473}
474
475#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
477pub struct SourceNormalizerResourceSessionStatus {
478 pub state: SourceNormalizerResourceSessionState,
479 #[serde(default)]
480 pub info: Option<SourceNormalizerResourceSessionInfo>,
481 #[serde(default)]
482 pub message: Option<String>,
483 #[serde(default)]
484 pub disk_bytes_used: Option<u64>,
485}
486
487#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
489pub struct SourceNormalizerResourceSessionWaitStatus {
490 pub updated: bool,
491}
492
493#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
495pub enum SourceNormalizerReadPacketStatus {
496 Packet,
497 NeedMoreData,
498 EndOfStream,
499}
500
501#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
503pub struct SourceNormalizerPacket {
504 pub pts_us: Option<i64>,
505 pub dts_us: Option<i64>,
506 pub duration_us: Option<i64>,
507 pub stream_index: u32,
508 #[serde(default)]
509 pub media_kind: SourceNormalizerPacketMediaKind,
510 pub key_frame: bool,
511 pub discontinuity: bool,
512 #[serde(default)]
513 pub sample_rate: Option<u32>,
514 #[serde(default)]
515 pub channels: Option<u16>,
516 #[serde(default)]
517 pub channel_layout: Option<String>,
518 #[serde(default)]
519 pub sample_format: Option<String>,
520 #[serde(default)]
521 pub frame_count: Option<u32>,
522 #[serde(default)]
523 pub end_of_stream: bool,
524}
525
526#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
528pub struct SourceNormalizerReadPacketMetadata {
529 pub status: SourceNormalizerReadPacketStatus,
530 #[serde(default)]
531 pub packet: Option<SourceNormalizerPacket>,
532 #[serde(default)]
533 pub message: Option<String>,
534}
535
536impl SourceNormalizerReadPacketMetadata {
537 pub fn packet(packet: SourceNormalizerPacket) -> Self {
538 Self {
539 status: SourceNormalizerReadPacketStatus::Packet,
540 packet: Some(packet),
541 message: None,
542 }
543 }
544
545 pub fn need_more_data(message: Option<String>) -> Self {
546 Self {
547 status: SourceNormalizerReadPacketStatus::NeedMoreData,
548 packet: None,
549 message,
550 }
551 }
552
553 pub fn end_of_stream() -> Self {
554 Self {
555 status: SourceNormalizerReadPacketStatus::EndOfStream,
556 packet: None,
557 message: None,
558 }
559 }
560}
561
562#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
564pub struct SourceNormalizerPacketSeek {
565 pub position_millis: u64,
566 #[serde(default)]
567 pub exact: bool,
568}
569
570#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
576pub struct SourceNormalizerOperationStatus {
577 pub completed: bool,
578 #[serde(default)]
579 pub message: Option<String>,
580}
581
582#[derive(Debug, Error, Clone, PartialEq, Eq, Serialize, Deserialize)]
584pub enum SourceNormalizerError {
585 #[error("unsupported runtime profile: {profile}")]
586 UnsupportedRuntimeProfile { profile: String },
587 #[error("invalid source normalizer input: {message}")]
588 InvalidInput { message: String },
589 #[error("source normalizer payload codec error: {message}")]
590 PayloadCodec { message: String },
591 #[error("source normalizer configuration error: {message}")]
592 Configuration { message: String },
593 #[error("source normalizer ABI violation: {message}")]
594 AbiViolation { message: String },
595 #[error("source normalizer session is not configured")]
596 NotConfigured,
597 #[error("source normalizer operation is unsupported: {operation}")]
598 UnsupportedOperation { operation: String },
599 #[error("source normalizer timeout: {message}")]
600 Timeout { message: String },
601 #[error("source normalizer resource exhausted: {message}")]
602 ResourceExhausted { message: String },
603 #[error("source normalizer internal error: {message}")]
604 Internal { message: String },
605}
606
607impl SourceNormalizerError {
608 pub fn payload_codec(message: impl Into<String>) -> Self {
609 Self::PayloadCodec {
610 message: message.into(),
611 }
612 }
613
614 pub fn configuration(message: impl Into<String>) -> Self {
615 Self::Configuration {
616 message: message.into(),
617 }
618 }
619
620 pub fn abi_violation(message: impl Into<String>) -> Self {
621 Self::AbiViolation {
622 message: message.into(),
623 }
624 }
625
626 pub fn invalid_input(message: impl Into<String>) -> Self {
627 Self::InvalidInput {
628 message: message.into(),
629 }
630 }
631
632 pub fn unsupported_operation(operation: impl Into<String>) -> Self {
633 Self::UnsupportedOperation {
634 operation: operation.into(),
635 }
636 }
637
638 pub fn internal(message: impl Into<String>) -> Self {
639 Self::Internal {
640 message: message.into(),
641 }
642 }
643
644 pub fn resource_exhausted(message: impl Into<String>) -> Self {
645 Self::ResourceExhausted {
646 message: message.into(),
647 }
648 }
649}
650
651#[doc(hidden)]
657pub fn validate_source_normalizer_plugin_input(
658 input: &str,
659 headers: &[(String, String)],
660) -> Result<(), SourceNormalizerError> {
661 if !headers.is_empty() {
662 return Err(SourceNormalizerError::invalid_input(
663 "Plugin sessions do not receive HTTP headers",
664 ));
665 }
666 if is_windows_drive_absolute_path(input) {
667 return Ok(());
668 }
669 if let Ok(url) = Url::parse(input)
670 && (!url.username().is_empty()
671 || url.password().is_some()
672 || url.query().is_some()
673 || url.fragment().is_some())
674 {
675 return Err(SourceNormalizerError::invalid_input(
676 "Plugin sessions do not receive URL credentials, query strings, or fragments",
677 ));
678 }
679 Ok(())
680}
681
682fn is_windows_drive_absolute_path(input: &str) -> bool {
683 let bytes = input.as_bytes();
684 bytes.len() >= 3
685 && bytes[0].is_ascii_alphabetic()
686 && bytes[1] == b':'
687 && matches!(bytes[2], b'/' | b'\\')
688}
689
690pub trait SourceNormalizerPacketPluginFactory: Send + Sync {
692 fn name(&self) -> &str;
693
694 fn packet_capabilities(&self) -> SourceNormalizerPacketCapabilities;
695
696 fn open_packet_session(
697 &self,
698 config: &SourceNormalizerPacketSessionConfig,
699 ) -> Result<Box<dyn SourceNormalizerPacketSession>, SourceNormalizerError>;
700}
701
702pub trait SourceNormalizerResourcePluginFactory: Send + Sync {
704 fn name(&self) -> &str;
705
706 fn resource_capabilities(&self) -> SourceNormalizerResourceCapabilities;
707
708 fn open_resource_session(
709 &self,
710 config: &SourceNormalizerResourceSessionConfig,
711 ) -> Result<Box<dyn SourceNormalizerResourceSession>, SourceNormalizerError>;
712}
713
714pub struct SourceNormalizerPacketLease<'a> {
716 pub metadata: SourceNormalizerReadPacketMetadata,
717 pub data: &'a [u8],
718 pub handle: usize,
719}
720
721impl std::fmt::Debug for SourceNormalizerPacketLease<'_> {
722 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
723 f.debug_struct("SourceNormalizerPacketLease")
724 .field("metadata", &self.metadata)
725 .field("data_len", &self.data.len())
726 .field("handle", &self.handle)
727 .finish()
728 }
729}
730
731pub trait SourceNormalizerPacketSession: Send {
733 fn stream_info(&self) -> SourceNormalizerPacketStreamInfo;
734
735 fn read_packet(&mut self) -> Result<SourceNormalizerPacketLease<'_>, SourceNormalizerError>;
736
737 fn release_packet(&mut self, packet_handle: usize) -> Result<(), SourceNormalizerError>;
738
739 fn seek(
740 &mut self,
741 seek: &SourceNormalizerPacketSeek,
742 ) -> Result<SourceNormalizerOperationStatus, SourceNormalizerError>;
743
744 fn flush(&mut self) -> Result<SourceNormalizerOperationStatus, SourceNormalizerError>;
745
746 fn close(&mut self) -> Result<(), SourceNormalizerError>;
747}
748
749pub trait SourceNormalizerResourceSession: Send {
751 fn session_info(&self) -> SourceNormalizerResourceSessionInfo;
752
753 fn poll(&mut self) -> Result<SourceNormalizerResourceSessionStatus, SourceNormalizerError>;
754
755 fn wait_for_update(
756 &mut self,
757 timeout: Duration,
758 ) -> Result<SourceNormalizerResourceSessionWaitStatus, SourceNormalizerError>;
759
760 fn cancel(&mut self) -> Result<SourceNormalizerOperationStatus, SourceNormalizerError>;
761
762 fn close(&mut self) -> Result<(), SourceNormalizerError>;
763}
764
765#[cfg(test)]
766mod tests {
767 use super::{
768 SourceNormalizerOutputRoute, SourceNormalizerPacket, SourceNormalizerPacketCapabilities,
769 SourceNormalizerPacketMediaKind, SourceNormalizerPacketSessionRequirements,
770 SourceNormalizerPacketTrackInfo, SourceNormalizerReadPacketMetadata,
771 SourceNormalizerReadPacketStatus, SourceNormalizerRequiredCapabilities,
772 SourceNormalizerResourceCachePolicy, SourceNormalizerResourceCapabilities,
773 SourceNormalizerResourceSessionRequirements, SourceNormalizerSessionCapabilities,
774 SourceNormalizerSessionRequirements,
775 };
776 use crate::{DecoderBitstreamFormat, NativeFrameColorMetadata, NativeFrameHdrMetadata};
777
778 #[test]
779 fn source_normalizer_resource_wait_status_round_trips_through_json() {
780 let status = super::SourceNormalizerResourceSessionWaitStatus { updated: true };
781 let json = serde_json::to_string(&status).expect("serialize wait status");
782 assert_eq!(json, r#"{"updated":true}"#);
783 let decoded: super::SourceNormalizerResourceSessionWaitStatus =
784 serde_json::from_str(&json).expect("decode wait status");
785 assert_eq!(decoded, status);
786 }
787
788 #[test]
789 fn source_normalizer_resource_capabilities_round_trip_through_json() {
790 let capabilities = SourceNormalizerResourceCapabilities {
791 supported_runtime_profiles: vec!["local-stream".to_owned()],
792 supported_output_routes: vec![
793 SourceNormalizerOutputRoute::Fmp4LocalStream,
794 SourceNormalizerOutputRoute::HlsShortWindow,
795 ],
796 max_level: Default::default(),
797 content_types: vec![
798 "video/mp4".to_owned(),
799 "application/vnd.apple.mpegurl".to_owned(),
800 ],
801 supports_growing_resources: true,
802 supports_range_reads: true,
803 supports_cancel: true,
804 required_capabilities: SourceNormalizerRequiredCapabilities::default(),
805 cache_policy: SourceNormalizerResourceCachePolicy::default(),
806 max_sessions: Some(2),
807 };
808
809 let encoded = serde_json::to_string(&capabilities).expect("serialize capabilities");
810 let decoded: SourceNormalizerResourceCapabilities =
811 serde_json::from_str(&encoded).expect("deserialize capabilities");
812
813 assert_eq!(decoded, capabilities);
814 assert!(decoded.supports_runtime_profile("LOCAL-STREAM"));
815 assert!(decoded.supports_output_route(SourceNormalizerOutputRoute::Fmp4LocalStream));
816 assert_eq!(
817 SourceNormalizerOutputRoute::HlsShortWindow.wire_name(),
818 "hlsShortWindow"
819 );
820 }
821
822 #[test]
823 fn source_normalizer_packet_metadata_round_trips_through_json() {
824 let metadata = SourceNormalizerReadPacketMetadata::packet(SourceNormalizerPacket {
825 pts_us: Some(33_000),
826 dts_us: Some(30_000),
827 duration_us: Some(33_333),
828 stream_index: 1,
829 media_kind: SourceNormalizerPacketMediaKind::Video,
830 key_frame: true,
831 discontinuity: false,
832 sample_rate: None,
833 channels: None,
834 channel_layout: None,
835 sample_format: None,
836 frame_count: None,
837 end_of_stream: false,
838 });
839
840 let encoded = serde_json::to_string(&metadata).expect("serialize packet metadata");
841 let decoded: SourceNormalizerReadPacketMetadata =
842 serde_json::from_str(&encoded).expect("deserialize packet metadata");
843
844 assert_eq!(decoded, metadata);
845 assert_eq!(decoded.status, SourceNormalizerReadPacketStatus::Packet);
846 }
847
848 #[test]
849 fn source_normalizer_packet_track_info_round_trips_through_json() {
850 let track = SourceNormalizerPacketTrackInfo {
851 stream_index: 0,
852 media_kind: SourceNormalizerPacketMediaKind::Video,
853 codec: "H264".to_owned(),
854 extradata: vec![1, 2, 3],
855 bitstream_format: Some(DecoderBitstreamFormat::Avcc),
856 width: Some(960),
857 height: Some(432),
858 coded_width: Some(960),
859 coded_height: Some(432),
860 reorder_depth: Some(4),
861 sample_rate: None,
862 channels: None,
863 channel_layout: None,
864 codec_delay_samples: None,
865 priming_samples: None,
866 trailing_padding_samples: None,
867 seek_preroll_samples: None,
868 color: Some(NativeFrameColorMetadata {
869 primaries: Some("bt2020".to_owned()),
870 transfer: Some("smpte2084".to_owned()),
871 matrix: Some("bt2020-ncl".to_owned()),
872 range: Some("limited".to_owned()),
873 bit_depth: Some(10),
874 }),
875 hdr: Some(NativeFrameHdrMetadata {
876 kind: "hdr10".to_owned(),
877 mastering_display: None,
878 content_light: None,
879 dolby_vision: None,
880 }),
881 frame_rate: Some(30.0),
882 time_base_num: Some(1),
883 time_base_den: Some(90_000),
884 };
885
886 let encoded = serde_json::to_string(&track).expect("serialize track");
887 let decoded: SourceNormalizerPacketTrackInfo =
888 serde_json::from_str(&encoded).expect("deserialize track");
889
890 assert_eq!(decoded, track);
891 assert_eq!(
892 decoded.color.as_ref().and_then(|color| color.bit_depth),
893 Some(10)
894 );
895 assert_eq!(
896 decoded.hdr.as_ref().map(|hdr| hdr.kind.as_str()),
897 Some("hdr10")
898 );
899 }
900
901 #[test]
902 fn source_normalizer_audio_packet_track_info_round_trips_through_json() {
903 let track = SourceNormalizerPacketTrackInfo {
904 stream_index: 1,
905 media_kind: SourceNormalizerPacketMediaKind::Audio,
906 codec: "AAC".to_owned(),
907 extradata: vec![0x12, 0x10],
908 bitstream_format: Some(DecoderBitstreamFormat::Unknown("AAC".to_owned())),
909 width: None,
910 height: None,
911 coded_width: None,
912 coded_height: None,
913 reorder_depth: None,
914 sample_rate: Some(48_000),
915 channels: Some(2),
916 channel_layout: Some("stereo".to_owned()),
917 codec_delay_samples: Some(0),
918 priming_samples: Some(2_112),
919 trailing_padding_samples: Some(512),
920 seek_preroll_samples: Some(1_024),
921 color: None,
922 hdr: None,
923 frame_rate: None,
924 time_base_num: Some(1),
925 time_base_den: Some(48_000),
926 };
927
928 let encoded = serde_json::to_string(&track).expect("serialize audio track");
929 let decoded: SourceNormalizerPacketTrackInfo =
930 serde_json::from_str(&encoded).expect("deserialize audio track");
931
932 assert_eq!(decoded, track);
933 }
934
935 #[test]
936 fn source_normalizer_packet_capabilities_support_case_insensitive_codecs() {
937 let capabilities = SourceNormalizerPacketCapabilities {
938 supported_runtime_profiles: vec!["diagnostic-packet".to_owned()],
939 max_level: Default::default(),
940 media_kinds: vec![SourceNormalizerPacketMediaKind::Video],
941 codecs: vec!["H264".to_owned()],
942 bitstream_formats: vec![DecoderBitstreamFormat::Avcc],
943 supports_seek: true,
944 supports_flush: true,
945 required_capabilities: SourceNormalizerRequiredCapabilities::default(),
946 max_sessions: Some(1),
947 };
948
949 assert!(capabilities.supports_runtime_profile("DIAGNOSTIC-PACKET"));
950 assert!(capabilities.supports_codec("h264"));
951 assert!(capabilities.supports_media_kind(SourceNormalizerPacketMediaKind::Video));
952 assert!(capabilities.supports_bitstream_format(&DecoderBitstreamFormat::Avcc));
953 }
954
955 #[test]
956 fn source_normalizer_packet_requirements_report_missing_capabilities() {
957 let requirements = SourceNormalizerPacketSessionRequirements {
958 require_seek: true,
959 bitstream_format: Some(DecoderBitstreamFormat::Avcc),
960 ..SourceNormalizerPacketSessionRequirements::native_video("native-frame-vod", "h264")
961 };
962 let capabilities = SourceNormalizerPacketCapabilities {
963 supported_runtime_profiles: vec!["other".to_owned()],
964 media_kinds: vec![SourceNormalizerPacketMediaKind::Audio],
965 codecs: vec!["hevc".to_owned()],
966 bitstream_formats: vec![DecoderBitstreamFormat::AnnexB],
967 supports_seek: false,
968 supports_flush: false,
969 ..Default::default()
970 };
971
972 let missing = requirements.missing_capabilities(&capabilities);
973
974 assert!(
975 missing
976 .iter()
977 .any(|item| item == "runtime profile native-frame-vod")
978 );
979 assert!(missing.iter().any(|item| item == "packet media kind Video"));
980 assert!(missing.iter().any(|item| item == "packet codec h264"));
981 assert!(
982 missing
983 .iter()
984 .any(|item| item == "packet bitstream format Avcc")
985 );
986 assert!(missing.iter().any(|item| item == "packet seek support"));
987 assert!(missing.iter().any(|item| item == "packet flush support"));
988 assert!(
989 missing
990 .iter()
991 .any(|item| item == "outstanding lease cleanup")
992 );
993 }
994
995 #[test]
996 fn source_normalizer_packet_requirements_treat_empty_profile_as_auto_detected() {
997 let requirements = SourceNormalizerPacketSessionRequirements {
998 runtime_profile: String::new(),
999 media_kind: Some(SourceNormalizerPacketMediaKind::Video),
1000 codec: None,
1001 bitstream_format: None,
1002 require_seek: false,
1003 require_flush: true,
1004 require_lease_cleanup: true,
1005 };
1006 let capabilities = SourceNormalizerPacketCapabilities {
1007 supported_runtime_profiles: vec!["native-frame-vod".to_owned()],
1008 media_kinds: vec![SourceNormalizerPacketMediaKind::Video],
1009 supports_flush: true,
1010 ..Default::default()
1011 };
1012
1013 assert!(requirements.missing_capabilities(&capabilities).is_empty());
1014 }
1015
1016 #[test]
1017 fn source_normalizer_session_requirements_match_route_kind() {
1018 let requirements = SourceNormalizerSessionRequirements::Resource(
1019 SourceNormalizerResourceSessionRequirements {
1020 runtime_profile: "vod-resource".to_owned(),
1021 output_route: SourceNormalizerOutputRoute::Fmp4LocalStream,
1022 content_type: Some("video/mp4".to_owned()),
1023 require_growing_resources: true,
1024 require_range_reads: true,
1025 require_cancel: true,
1026 },
1027 );
1028 let packet_capabilities = SourceNormalizerPacketCapabilities::default();
1029 let missing = requirements.missing_capabilities(
1030 &SourceNormalizerSessionCapabilities::Packet(&packet_capabilities),
1031 );
1032 assert_eq!(missing, vec!["resource output route".to_owned()]);
1033
1034 let resource_capabilities = SourceNormalizerResourceCapabilities {
1035 supported_runtime_profiles: vec!["vod-resource".to_owned()],
1036 supported_output_routes: vec![SourceNormalizerOutputRoute::Fmp4LocalStream],
1037 content_types: vec!["video/mp4".to_owned()],
1038 supports_growing_resources: true,
1039 supports_range_reads: true,
1040 supports_cancel: true,
1041 ..Default::default()
1042 };
1043 let missing = requirements.missing_capabilities(
1044 &SourceNormalizerSessionCapabilities::Resource(&resource_capabilities),
1045 );
1046 assert!(missing.is_empty());
1047 }
1048}