use std::time::Duration;
use rskit_errors::AppResult;
use rskit_storage::{FileSink, FileSource};
use tokio_util::sync::CancellationToken;
use crate::{
executor::MediaExecutor,
filter::Filter,
ops::{
ConcatOp, CropRegion, FilterConfig, FlipDirection, InterpolateConfig, MediaOp, MixAudioOp,
OverlayConfig, OverlayOp, OverlayPosition, PadOp, ReplaceAudioOp, ResizeMode, ResizeOp,
Rotation, SceneDetectConfig, SubtitleConfig, ThumbnailConfig, Transition, UpscaleConfig,
},
output::OutputConfig,
spatial::Resolution,
subtitle::SubtitleTrack,
time::{Segment, TimeRange, Timestamp},
types::TrackKind,
};
pub struct MediaPipeline {
source: FileSource,
ops: Vec<MediaOp>,
sink: Option<FileSink>,
}
impl MediaPipeline {
pub fn from(source: &FileSource) -> Self {
Self {
source: source.clone(),
ops: Vec::new(),
sink: None,
}
}
#[must_use]
pub fn extract(mut self, range: TimeRange) -> Self {
self.ops.push(MediaOp::Extract(range));
self
}
#[must_use]
pub fn extract_many(mut self, segments: Vec<Segment>) -> Self {
self.ops.push(MediaOp::ExtractMany(segments));
self
}
#[must_use]
pub fn resize(mut self, resolution: Resolution, mode: ResizeMode) -> Self {
self.ops
.push(MediaOp::Resize(ResizeOp { resolution, mode }));
self
}
#[must_use]
pub fn crop(mut self, region: CropRegion) -> Self {
self.ops.push(MediaOp::Crop(region));
self
}
#[must_use]
pub fn rotate(mut self, rotation: Rotation) -> Self {
self.ops.push(MediaOp::Rotate(rotation));
self
}
#[must_use]
pub fn flip(mut self, direction: FlipDirection) -> Self {
self.ops.push(MediaOp::Flip(direction));
self
}
#[must_use]
pub fn pad(mut self, width: u32, height: u32, color: &str) -> Self {
self.ops.push(MediaOp::Pad(PadOp {
width,
height,
color: color.to_string(),
}));
self
}
#[must_use]
pub fn speed(mut self, factor: f64) -> Self {
self.ops.push(MediaOp::Speed(factor));
self
}
#[must_use]
pub fn reverse(mut self) -> Self {
self.ops.push(MediaOp::Reverse);
self
}
#[must_use]
pub fn volume(mut self, factor: f64) -> Self {
self.ops.push(MediaOp::Volume(factor));
self
}
#[must_use]
pub fn normalize_audio(mut self) -> Self {
self.ops.push(MediaOp::NormalizeAudio);
self
}
#[must_use]
pub fn fade_in(mut self, duration: Duration) -> Self {
self.ops.push(MediaOp::FadeIn(duration));
self
}
#[must_use]
pub fn fade_out(mut self, duration: Duration) -> Self {
self.ops.push(MediaOp::FadeOut(duration));
self
}
#[must_use]
pub fn strip_audio(mut self) -> Self {
self.ops.push(MediaOp::StripAudio);
self
}
#[must_use]
pub fn strip_video(mut self) -> Self {
self.ops.push(MediaOp::StripVideo);
self
}
#[must_use]
pub fn filter(mut self, filter: Filter) -> Self {
self.ops.push(MediaOp::Filter(filter));
self
}
#[must_use]
pub fn overlay(mut self, source: &FileSource, position: OverlayPosition, opacity: f32) -> Self {
self.ops.push(MediaOp::Overlay(OverlayOp {
source: source.clone(),
position,
opacity,
time_range: None,
scale: None,
}));
self
}
#[must_use]
pub fn concat(mut self, source: &FileSource) -> Self {
self.ops.push(MediaOp::Concat(ConcatOp {
source: source.clone(),
transition: None,
}));
self
}
#[must_use]
pub fn concat_with_transition(mut self, source: &FileSource, transition: Transition) -> Self {
self.ops.push(MediaOp::Concat(ConcatOp {
source: source.clone(),
transition: Some(transition),
}));
self
}
#[must_use]
pub fn replace_audio(mut self, audio: &FileSource) -> Self {
self.ops.push(MediaOp::ReplaceAudio(ReplaceAudioOp {
audio_source: audio.clone(),
offset: None,
}));
self
}
#[must_use]
pub fn mix_audio(mut self, audio: &FileSource, volume: f64) -> Self {
self.ops.push(MediaOp::MixAudio(MixAudioOp {
audio_source: audio.clone(),
volume,
offset: None,
}));
self
}
#[must_use]
pub fn burn_subtitles(mut self, subs: SubtitleTrack) -> Self {
self.ops.push(MediaOp::BurnSubtitles(subs));
self
}
#[must_use]
pub fn apply_filter(mut self, config: FilterConfig) -> Self {
self.ops.push(MediaOp::ApplyFilter(config));
self
}
#[must_use]
pub fn add_overlay(mut self, config: OverlayConfig) -> Self {
self.ops.push(MediaOp::AddOverlay(config));
self
}
#[must_use]
pub fn generate_thumbnail(mut self, config: ThumbnailConfig) -> Self {
self.ops.push(MediaOp::GenerateThumbnail(config));
self
}
#[must_use]
pub fn detect_scenes(mut self, config: SceneDetectConfig) -> Self {
self.ops.push(MediaOp::DetectScenes(config));
self
}
#[must_use]
pub fn add_subtitles(mut self, config: SubtitleConfig) -> Self {
self.ops.push(MediaOp::AddSubtitles(config));
self
}
#[must_use]
pub fn upscale(mut self, config: UpscaleConfig) -> Self {
self.ops.push(MediaOp::Upscale(config));
self
}
#[must_use]
pub fn interpolate(mut self, config: InterpolateConfig) -> Self {
self.ops.push(MediaOp::Interpolate(config));
self
}
#[must_use]
pub fn select_tracks(mut self, indices: Vec<usize>) -> Self {
self.ops.push(MediaOp::SelectTracks(indices));
self
}
#[must_use]
pub fn select_tracks_by_kind(mut self, kinds: Vec<TrackKind>) -> Self {
self.ops.push(MediaOp::SelectTracksByKind(kinds));
self
}
#[must_use]
pub fn transcode(mut self, config: OutputConfig) -> Self {
self.ops.push(MediaOp::Transcode(config));
self
}
#[must_use]
pub fn stream_to(mut self, config: OutputConfig) -> Self {
self.ops.push(MediaOp::Transcode(config));
self
}
#[must_use]
pub fn output_to(mut self, sink: FileSink) -> Self {
self.sink = Some(sink);
self
}
pub fn validate(&self, executor: &dyn MediaExecutor) -> AppResult<()> {
for op in &self.ops {
if !executor.supports(op) {
return Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::InvalidInput,
format!("executor does not support operation: {op:?}"),
));
}
}
let mut has_strip_audio = false;
let mut has_strip_video = false;
for op in &self.ops {
if matches!(op, MediaOp::StripAudio) {
has_strip_audio = true;
continue;
}
if matches!(op, MediaOp::StripVideo) {
has_strip_video = true;
continue;
}
if has_strip_audio && op.requires_audio_track() {
return Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::InvalidInput,
"cannot apply audio operation after StripAudio",
));
}
if has_strip_video && op.requires_video_track() {
return Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::InvalidInput,
"cannot apply video operation after StripVideo",
));
}
}
Ok(())
}
pub async fn execute(self, executor: &dyn MediaExecutor) -> AppResult<FileSource> {
executor
.execute(&self.source, &self.ops, self.sink.as_ref())
.await
}
pub async fn execute_with_progress(
self,
executor: &dyn MediaExecutor,
on_progress: impl Fn(Progress) + Send + Sync + 'static,
) -> AppResult<FileSource> {
executor
.execute_with_progress(
&self.source,
&self.ops,
self.sink.as_ref(),
Box::new(on_progress),
)
.await
}
pub async fn execute_cancellable(
self,
executor: &dyn MediaExecutor,
on_progress: Option<Box<dyn Fn(Progress) + Send + Sync>>,
cancel: CancellationToken,
) -> AppResult<FileSource> {
executor
.execute_cancellable(
&self.source,
&self.ops,
self.sink.as_ref(),
on_progress,
cancel,
)
.await
}
pub fn operations(&self) -> &[MediaOp] {
&self.ops
}
pub fn estimated_duration(&self, source_duration: Duration) -> Duration {
let mut duration = source_duration;
for op in &self.ops {
match op {
MediaOp::Extract(range) => {
duration = range.duration();
}
MediaOp::ExtractMany(segments) => {
let total_us: u64 = segments
.iter()
.map(|s| s.range.duration().as_micros() as u64)
.sum();
duration = Duration::from_micros(total_us);
}
MediaOp::Concat(_) => {
duration = duration.saturating_mul(2);
}
MediaOp::Speed(factor) if *factor > 0.0 => {
let us = (duration.as_micros() as f64 / factor) as u64;
duration = Duration::from_micros(us);
}
_ => {}
}
}
duration
}
}
#[derive(Debug, Clone)]
pub struct Progress {
pub position: Option<Timestamp>,
pub total: Option<Duration>,
pub percent: Option<f32>,
pub speed: Option<f64>,
pub output_size: Option<u64>,
pub eta: Option<Duration>,
}