Skip to content

Commit aa1be0a

Browse files
committed
plugin transcode
1 parent 3427d01 commit aa1be0a

10 files changed

Lines changed: 391 additions & 146 deletions

File tree

Cargo.lock

Lines changed: 2 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ sha256 = "1.5.0"
7070
axum-extra = { version = "0.9.2", features = ["query"] }
7171
http = "1.1.0"
7272
extism = "1.10.0"
73-
rs-plugin-common-interfaces = { version = "0.19.0", features = ["rusqlite",] }
73+
rs-plugin-common-interfaces = { version = "0.20.3", features = ["rusqlite",] }
7474
async-recursion = "1.1.0"
7575
async-compression = { version = "0.4.6", features = ["tokio"] }
7676
youtube_dl = { version = "0.10.0", features = ["tokio", "downloader-rustls-tls"] }

src/domain/media.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
11
use std::str::FromStr;
22

33
use nanoid::nanoid;
4-
use rs_plugin_common_interfaces::{request::{RsCookie, RsRequest, RsRequestStatus}, url::RsLink};
4+
use rs_plugin_common_interfaces::{request::{RsCookie, RsRequest, RsRequestStatus}, url::RsLink, video::VideoConvertRequest};
55
use serde::{Deserialize, Serialize};
66
use serde_json::Value;
77
use strum_macros::EnumString;
88

99

10-
use crate::{plugins::sources::SourceRead, tools::video_tools::VideoConvertRequest};
10+
use crate::plugins::sources::SourceRead;
1111

1212
use super::{progress::RsProgress, ElementAction};
1313

src/model/medias.rs

Lines changed: 115 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -13,16 +13,16 @@ use mime_guess::get_mime_extensions_str;
1313
use nanoid::nanoid;
1414
use query_external_ip::SourceError;
1515
use regex::Regex;
16-
use rs_plugin_common_interfaces::{request::{RsRequest, RsRequestStatus}, url::{RsLink, RsLinkType}, PluginType, RsCookie};
16+
use rs_plugin_common_interfaces::{request::{RsRequest, RsRequestStatus}, url::{RsLink, RsLinkType}, video::{RsVideoTranscodeJob, RsVideoTranscodeJobPluginRequest, RsVideoTranscodeStatus, VideoConvertRequest, VideoOverlayType}, PluginType, RsCookie};
1717
use rusqlite::{types::{FromSql, FromSqlError, FromSqlResult, ToSqlOutput, ValueRef}, ToSql};
1818
use serde::{Deserialize, Serialize};
1919
use strum_macros::EnumString;
20-
use tokio::{fs::File, io::{copy, AsyncRead, AsyncReadExt, AsyncWriteExt}, sync::mpsc};
20+
use tokio::{fs::File, io::{copy, AsyncRead, AsyncReadExt, AsyncWriteExt}, sync::mpsc, time::sleep, time::Duration};
2121
use tokio_stream::StreamExt;
2222
use tokio_util::{compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt}, io::{ReaderStream, StreamReader, SyncIoBridge}};
2323
use zip::ZipWriter;
2424

25-
use crate::{domain::{deleted::RsDeleted, library::LibraryType, media::{self, ConvertMessage, ConvertProgress, RsGpsPosition, DEFAULT_MIME}, plugin, MediaElement}, error::RsError, model::store::sql::SqlOrder, plugins::sources::{path_provider::PathProvider, Source}, routes::infos, tools::{file_tools::{filename_from_path, remove_extension}, image_tools::{convert_image_reader, image_infos, IMAGES_MIME_FULL_BROWSER_SUPPORT}, recognition, video_tools::{VideoCommandBuilder, VideoConvertRequest, VideoOverlayPosition}}};
25+
use crate::{domain::{deleted::RsDeleted, library::LibraryType, media::{self, ConvertMessage, ConvertProgress, RsGpsPosition, DEFAULT_MIME}, plugin, MediaElement}, error::RsError, model::store::sql::SqlOrder, plugins::sources::{path_provider::PathProvider, Source}, routes::infos, tools::{file_tools::{filename_from_path, remove_extension}, image_tools::{convert_image_reader, image_infos, IMAGES_MIME_FULL_BROWSER_SUPPORT}, recognition, video_tools::{VideoCommandBuilder}}};
2626

2727
use crate::{domain::{library::LibraryRole, media::{FileType, GroupMediaDownload, Media, MediaDownloadUrl, MediaForAdd, MediaForInsert, MediaForUpdate, MediaItemReference, MediaWithAction, MediasMessage, ProgressMessage}, progress::{RsProgress, RsProgressType}, ElementAction}, error::RsResult, plugins::{get_plugin_fodler, sources::{async_reader_progress::ProgressReader, error::SourcesError, AsyncReadPinBox, FileStreamResult, SourceRead}}, routes::mw_range::RangeDefinition, server::get_server_port, tools::{auth::{sign_local, ClaimsLocal}, file_tools::{file_type_from_mime, get_extension_from_mime}, image_tools::{self, resize_image_reader, ImageSize}, log::{log_error, log_info, LogServiceType}, prediction::{predict_net, preload_model, PredictionTagResult}, video_tools::{self, probe_video, VideoTime}}};
2828

@@ -1050,16 +1050,16 @@ impl ModelController {
10501050
let uri = if let Some(local_path) = local_path {
10511051
local_path.to_str().unwrap().to_string()
10521052
} else {
1053-
ModelController::get_temporary_local_read_url(library_id, media_id).await?
1053+
ModelController::get_temporary_local_read_url(library_id, media_id, None).await?
10541054
};
10551055
let thumb = video_tools::thumb_video(&uri, time).await?;
10561056
let mut cursor = std::io::Cursor::new(thumb);
10571057
let thumb = resize_image_reader(Box::pin(cursor), 512, format, quality, false).await?;
10581058
Ok(thumb)
10591059
}
10601060

1061-
pub async fn get_temporary_local_read_url(library_id: &str, media_id: &str) -> Result<String> {
1062-
let exp = ClaimsLocal::generate_seconds(240);
1061+
pub async fn get_temporary_local_read_url(library_id: &str, media_id: &str, delay: Option<u64>) -> Result<String> {
1062+
let exp = ClaimsLocal::generate_seconds(delay.unwrap_or(240));
10631063
let claims = ClaimsLocal {
10641064
cr: "service::get_video_thumb".to_string(),
10651065
kind: crate::tools::auth::ClaimsLocalType::File(library_id.to_string(), media_id.to_string()),
@@ -1127,20 +1127,124 @@ impl ModelController {
11271127
}
11281128
}
11291129

1130-
pub async fn convert(&self, library_id: &str, media_id: &str, request: VideoConvertRequest, requesting_user: &ConnectedUser) -> crate::Result<()> {
1130+
1131+
pub async fn convert(&self, library_id: &str, media_id: &str, request: VideoConvertRequest, plugin_id: Option<String>, requesting_user: &ConnectedUser) -> crate::Result<()> {
11311132
requesting_user.check_file_role(library_id, media_id, LibraryRole::Write)?;
11321133
let store = self.store.get_library_store(library_id)?;
11331134
let media = store.get_media(media_id, requesting_user.user_id().ok()).await?.ok_or(SourcesError::UnableToFindMedia(library_id.to_string(), media_id.to_string(), "convert".to_string()))?;
11341135

11351136

11361137
let filename = format!("{}.{}", remove_extension(&media.name), request.format);
1137-
let queue_element = VideoConvertQueueElement::new(library_id.to_string(), media_id.to_string(), filename, requesting_user.clone(), request);
1138+
let queue_element = VideoConvertQueueElement::new(library_id.to_string(), plugin_id.clone(), media_id.to_string(), filename.clone(), requesting_user.clone(), request.clone());
11381139
let message = ConvertMessage {
11391140
library: library_id.to_string(),
11401141
progress: queue_element.status.clone(),
11411142
};
11421143
self.send_convert_progress(message);
11431144

1145+
1146+
// If plugin convert is requested
1147+
if let Some(plugin_id) = &plugin_id {
1148+
let status = self.convert_submit_media(library_id, media_id, request.clone(), &plugin_id).await?;
1149+
let job_id = status.id.clone();
1150+
1151+
// Clone what you need for the background task
1152+
let plugin_id_owned = plugin_id.clone();
1153+
let mc_progress = self.clone(); // or however you get your progress sender
1154+
let request_clone = request.clone();
1155+
let lib_progress = library_id.to_string(); // adjust based on your actual type
1156+
let name_progress = filename.clone();
1157+
let media_progress: MediaForUpdate = media.into();
1158+
// Spawn background task
1159+
tokio::spawn(async move {
1160+
let mut poll_interval = Duration::from_secs(2);
1161+
1162+
loop {
1163+
sleep(poll_interval).await;
1164+
1165+
match mc_progress.convert_status(&job_id, &plugin_id_owned).await {
1166+
Ok(current_status) => {
1167+
let is_terminal = matches!(
1168+
current_status.status,
1169+
RsVideoTranscodeStatus::Completed
1170+
| RsVideoTranscodeStatus::Failed
1171+
| RsVideoTranscodeStatus::Canceled
1172+
);
1173+
1174+
let done = is_terminal;
1175+
let percent = current_status.progress;
1176+
1177+
let message = ConvertMessage {
1178+
library: lib_progress.clone(),
1179+
progress: ConvertProgress {
1180+
percent: percent.into(),
1181+
converted_id: if done { Some(job_id.clone()) } else { None },
1182+
filename: name_progress.clone(),
1183+
done,
1184+
id: request_clone.id.clone(),
1185+
request: Some(request_clone.clone()),
1186+
estimated_remaining_seconds: None,
1187+
},
1188+
};
1189+
mc_progress.send_convert_progress(message);
1190+
1191+
if matches!(current_status.status, RsVideoTranscodeStatus::Completed) {
1192+
// Handle all errors here instead of using ?
1193+
match mc_progress.convert_link(&job_id, &plugin_id_owned).await {
1194+
Ok(link) => {
1195+
let source = SourceRead::Request(link);
1196+
1197+
match source.into_reader(
1198+
Some(&lib_progress),
1199+
None,
1200+
None,
1201+
Some((mc_progress.clone(), &ConnectedUser::ServerAdmin)),
1202+
None
1203+
).await {
1204+
Ok(reader) => {
1205+
// Note: Can't use self here, need to use mc_progress
1206+
match mc_progress.add_library_file(
1207+
&lib_progress,
1208+
&name_progress,
1209+
Some(media_progress.clone()),
1210+
reader.stream,
1211+
&ConnectedUser::ServerAdmin // Or pass requesting_user as owned
1212+
).await {
1213+
Ok(media) => {
1214+
// Success - maybe send final progress update
1215+
tracing::info!("Successfully added media: {:?}", media);
1216+
}
1217+
Err(e) => {
1218+
tracing::error!("Failed to add library file: {:?}", e);
1219+
// Send error progress update
1220+
}
1221+
}
1222+
}
1223+
Err(e) => {
1224+
tracing::error!("Failed to create reader: {:?}", e);
1225+
}
1226+
}
1227+
}
1228+
Err(e) => {
1229+
tracing::error!("Failed to get convert link: {:?}", e);
1230+
}
1231+
}
1232+
}
1233+
1234+
if is_terminal {
1235+
break;
1236+
}
1237+
}
1238+
Err(e) => {
1239+
tracing::error!("Error polling status: {:?}", e);
1240+
break;
1241+
}
1242+
}
1243+
}
1244+
});
1245+
return Ok(());
1246+
}
1247+
11441248
let converting = self.convert_current.read().await;
11451249
let mut queue = self.convert_queue.write().await;
11461250
queue.push_back(queue_element);
@@ -1233,11 +1337,11 @@ impl ModelController {
12331337

12341338
if let Some(overlay) = &mut element.request.overlay {
12351339
match overlay.kind {
1236-
video_tools::VideoOverlayType::Watermark => {
1340+
VideoOverlayType::Watermark => {
12371341
let name = if overlay.path.is_empty() { ".watermark.png".to_owned() } else { format!(".watermark.{}.png", &overlay.path)};
12381342
overlay.path = local.get_full_path(&name).to_str().ok_or(Error::ServiceError("Convert".to_owned(), Some("Invalid watermark path".to_owned())))?.to_string();
12391343
},
1240-
video_tools::VideoOverlayType::File => todo!(),
1344+
VideoOverlayType::File => todo!(),
12411345
}
12421346

12431347
}
@@ -1301,7 +1405,7 @@ impl ModelController {
13011405
let uri = if let Some(local_path) = local_path {
13021406
local_path.to_str().unwrap().to_string()
13031407
} else {
1304-
ModelController::get_temporary_local_read_url(library_id, media_id).await?
1408+
ModelController::get_temporary_local_read_url(library_id, media_id, Some(240)).await?
13051409
};
13061410

13071411
let videos_infos = probe_video(&uri).await?;

src/model/mod.rs

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,10 +21,10 @@ pub mod player;
2121
use std::{collections::{HashMap, VecDeque}, io::Read, path::PathBuf, pin::Pin, sync::Arc, thread::JoinHandle};
2222
use futures::lock::Mutex;
2323
use nanoid::nanoid;
24-
use rs_plugin_common_interfaces::{ImageType, RsRequest};
24+
use rs_plugin_common_interfaces::{video::VideoConvertRequest, ImageType, RsRequest};
2525
use serde::{Deserialize, Serialize};
2626
use strum::IntoEnumIterator;
27-
use crate::{domain::{backup::BackupProcessStatus, library::{LibraryMessage, LibraryRole, ServerLibrary}, media::ConvertProgress, player::RsPlayerAvailable}, error::{RsError, RsResult}, plugins::{medias::{fanart::FanArtContext, imdb::ImdbContext, tmdb::TmdbContext, trakt::TraktContext}, sources::{error::SourcesError, local_provider_for_library, path_provider::PathProvider, AsyncReadPinBox, FileStreamResult, Source, SourceRead}, PluginManager}, tools::{clock::SECONDS_IN_HOUR, image_tools::{resize_image_reader, ImageSize}, log::log_info, scheduler::{self, ip::RefreshIpTask, refresh::RefreshTask, RsScheduler, RsTaskType}, video_tools::VideoConvertRequest}};
27+
use crate::{domain::{backup::BackupProcessStatus, library::{LibraryMessage, LibraryRole, ServerLibrary}, media::ConvertProgress, player::RsPlayerAvailable}, error::{RsError, RsResult}, plugins::{medias::{fanart::FanArtContext, imdb::ImdbContext, tmdb::TmdbContext, trakt::TraktContext}, sources::{error::SourcesError, local_provider_for_library, path_provider::PathProvider, AsyncReadPinBox, FileStreamResult, Source, SourceRead}, PluginManager}, tools::{clock::SECONDS_IN_HOUR, image_tools::{resize_image_reader, ImageSize}, log::log_info, scheduler::{self, ip::RefreshIpTask, refresh::RefreshTask, RsScheduler, RsTaskType}}};
2828

2929
use self::{medias::CRYPTO_HEADER_SIZE, store::SqliteStore, users::{ConnectedUser, ServerUser, UserRole}};
3030
use error::{Result, Error};
@@ -38,13 +38,15 @@ pub struct VideoConvertQueueElement {
3838
media: String,
3939
user: ConnectedUser,
4040
id: String,
41+
plugin_id: Option<String>,
4142
status: ConvertProgress
4243

4344
}
4445

4546
impl VideoConvertQueueElement {
46-
pub fn new(library: String, media: String, filename: String, user: ConnectedUser, request: VideoConvertRequest) -> VideoConvertQueueElement {
47+
pub fn new(library: String, plugin_id: Option<String>, media: String, filename: String, user: ConnectedUser, request: VideoConvertRequest) -> VideoConvertQueueElement {
4748
VideoConvertQueueElement {id: request.id.clone(),
49+
plugin_id,
4850
status: ConvertProgress { id: request.id.clone(), filename, converted_id: None, done: false, percent: 0f64, estimated_remaining_seconds: None, request: Some(request.clone()) },
4951
request, library, media, user,
5052
}
Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@ use crate::{domain::{backup::Backup, library::LibraryRole, plugin::{Plugin, Plug
1515

1616
use super::{error::{Error, Result}, users::{ConnectedUser, UserRole}, ModelController};
1717

18+
pub mod video_convert_plugin;
19+
1820
#[derive(Debug, Serialize, Deserialize, Clone, Default)]
1921
pub struct PluginQuery {
2022
pub kind: Option<PluginType>,
@@ -55,7 +57,7 @@ impl ModelController {
5557
Ok(plugins)
5658
}
5759

58-
async fn get_plugins_with_credential(&self, query: PluginQuery) -> Result<impl Iterator<Item = PluginWithCredential>> {
60+
pub async fn get_plugins_with_credential(&self, query: PluginQuery) -> Result<impl Iterator<Item = PluginWithCredential>> {
5961
let plugins = self.store.get_plugins(query).await?.into_iter();
6062
let credentials = self.store.get_credentials().await?;
6163
let iter = plugins.map(move |p| {
@@ -70,6 +72,14 @@ impl ModelController {
7072
let credential = self.store.get_plugin(&plugin_id).await?.ok_or(SourcesError::UnableToFindPlugin(plugin_id.to_string(), "get_plugin".to_string()))?;
7173
Ok(credential)
7274
}
75+
76+
pub async fn get_plugin_with_credential(&self, id: &str) -> Result<PluginWithCredential> {
77+
let plugin = self.store.get_plugin(id).await?.ok_or(SourcesError::UnableToFindPlugin(id.to_string(), "get_plugin_with_credential".to_string()))?;
78+
let credentials = self.store.get_credentials().await?;
79+
80+
let credential = credentials.iter().find(|c| Some(&c.id) == plugin.credential.as_ref()).cloned();
81+
Ok(PluginWithCredential { plugin: plugin, credential })
82+
}
7383

7484
pub async fn reload_plugins(&self, requesting_user: &ConnectedUser) -> RsResult<()> {
7585
requesting_user.check_role(&UserRole::Admin)?;

0 commit comments

Comments
 (0)