Skip to main content

player_plugin/
source_normalizer.rs

1use std::time::Duration;
2
3use serde::{Deserialize, Serialize};
4use thiserror::Error;
5use url::Url;
6
7use crate::{DecoderBitstreamFormat, NativeFrameColorMetadata, NativeFrameHdrMetadata};
8
9/// Normalization work level supported by a source normalizer plugin.
10#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default, Serialize, Deserialize)]
11pub enum SourceNormalizerNormalizeLevel {
12    /// Remux/copy normalization with optional bitstream filters.
13    #[default]
14    #[serde(alias = "remux_only", alias = "remux-only")]
15    RemuxOnly = 1,
16    /// Packet repair that still does not decode media into frames.
17    #[serde(alias = "packet_repair", alias = "packet-repair")]
18    PacketRepair = 2,
19}
20
21/// FFmpeg-like build features required by a source normalizer profile.
22#[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/// Capabilities advertised by a packet-stream source normalizer plugin.
37#[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    /// Returns whether this plugin advertises a runtime profile.
52    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    /// Returns whether this plugin advertises a codec.
59    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    /// Returns whether this plugin advertises packet output for a media kind.
66    pub fn supports_media_kind(&self, media_kind: SourceNormalizerPacketMediaKind) -> bool {
67        self.media_kinds.contains(&media_kind)
68    }
69
70    /// Returns whether this plugin advertises a packet bitstream format.
71    pub fn supports_bitstream_format(&self, format: &DecoderBitstreamFormat) -> bool {
72        self.bitstream_formats.contains(format)
73    }
74}
75
76/// Normalized output route produced by a source normalizer.
77#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
78#[serde(rename_all = "camelCase")]
79pub enum SourceNormalizerOutputRoute {
80    /// Disk-backed fragmented MP4 output intended to be exposed as a local stream.
81    Fmp4LocalStream,
82    /// Disk-backed short-window HLS output intended for nonstandard adaptive input.
83    HlsShortWindow,
84    /// Compressed packet stream intended for the SDK-controlled native frame lane.
85    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/// Resource session cache limits shared by plugin and platform hosts.
99#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
100pub struct SourceNormalizerResourceCachePolicy {
101    /// Maximum bytes read into memory per active session.
102    pub session_read_buffer_bytes: u64,
103    /// Maximum bytes used for manifest and metadata snapshots per session.
104    pub manifest_snapshot_bytes: u64,
105    /// Soft disk limit for one resource session.
106    pub session_disk_soft_cap_bytes: u64,
107    /// Soft disk limit for all normalized-resource sessions owned by a host.
108    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/// Capabilities advertised by a resource-output source normalizer plugin.
123#[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    /// Returns whether this plugin advertises a runtime profile.
139    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    /// Returns whether this plugin advertises an output route.
146    pub fn supports_output_route(&self, route: SourceNormalizerOutputRoute) -> bool {
147        self.supported_output_routes.contains(&route)
148    }
149}
150
151/// Capability requirements used when opening one source normalizer session.
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
153pub enum SourceNormalizerSessionRequirements {
154    Packet(SourceNormalizerPacketSessionRequirements),
155    Resource(SourceNormalizerResourceSessionRequirements),
156}
157
158impl SourceNormalizerSessionRequirements {
159    /// Returns missing capability names for this requirement.
160    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/// Borrowed source normalizer capabilities used for generic requirement matching.
184#[derive(Debug, Clone, Copy)]
185pub enum SourceNormalizerSessionCapabilities<'a> {
186    Packet(&'a SourceNormalizerPacketCapabilities),
187    Resource(&'a SourceNormalizerResourceCapabilities),
188}
189
190/// Packet-stream capability requirements.
191#[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    /// Builds packet-stream requirements for native-frame video decode.
210    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    /// Returns missing capability names for this requirement.
223    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/// Resource-output capability requirements.
262#[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    /// Returns missing capability names for this requirement.
278    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/// Packet stream media kind produced by a source normalizer.
314#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default, Serialize, Deserialize)]
315pub enum SourceNormalizerPacketMediaKind {
316    #[default]
317    Video,
318    Audio,
319    Subtitle,
320}
321
322/// Configuration used to open one packet-stream source normalizer session.
323#[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/// Configuration used to open one disk-backed resource source normalizer session.
338#[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/// Track metadata exposed by a packet-stream source normalizer.
356#[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/// Packet-stream metadata returned after opening a source normalizer session.
402#[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/// Disk-backed resource produced by a source normalizer session.
422#[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/// Resource-output metadata returned after opening a source normalizer session.
435#[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/// Resource-output worker state returned by `SourceNormalizerResourceSession::poll`.
464#[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/// Resource-output worker status returned by a source normalizer session.
476#[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/// Result returned after waiting for a resource-output session update.
488#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
489pub struct SourceNormalizerResourceSessionWaitStatus {
490    pub updated: bool,
491}
492
493/// Packet read status encoded in source normalizer packet metadata.
494#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
495pub enum SourceNormalizerReadPacketStatus {
496    Packet,
497    NeedMoreData,
498    EndOfStream,
499}
500
501/// Compressed packet metadata returned by a packet-stream source normalizer.
502#[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/// Metadata returned by `SourceNormalizerPacketSession::read_packet`.
527#[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/// Seek request passed to an active packet-stream source normalizer session.
563#[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/// Success payload used by source-normalizer session operations.
571///
572/// Resource cancellation may report `completed = false` after accepting the
573/// request while background work is still terminating. Callers observe the
574/// terminal state through `poll` or `wait_for_update`.
575#[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/// Error payload shared by source normalizer plugins and host-side adapters.
583#[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/// Validates the intentionally restricted locator surface exposed to plugins.
652///
653/// This shared host/adapter boundary rejects sensitive data instead of
654/// rewriting the locator, because query and fragment removal would change its
655/// identity. Local paths retain literal `?` and `#` characters.
656#[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
690/// Creates packet-stream source normalizer sessions for one plugin.
691pub 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
702/// Creates resource-output source normalizer sessions for one plugin.
703pub 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
714/// Borrowed packet returned by a packet-stream source normalizer.
715pub 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
731/// Stateful packet-stream source normalizer session.
732pub 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
749/// Stateful resource-output source normalizer session.
750pub 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}