Skip to content

Commit 2d704c5

Browse files
committed
Improve progress reporting and media source buffering
- Rename remaining ETA field in SSE/media convert payloads - Add progress estimation and retry handling for media imports - Increase stream and file I/O buffers for source providers
1 parent 894fe1f commit 2d704c5

9 files changed

Lines changed: 149 additions & 49 deletions

File tree

docs/SSE.md

Lines changed: 32 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -122,18 +122,41 @@ interface Relations {
122122
booksDetails?: Book[];
123123
}
124124

125+
type RsProgressType =
126+
| 'download'
127+
| 'transfert'
128+
| 'analysing'
129+
| 'finished'
130+
| { duplicate: string };
131+
132+
interface RsProgress {
133+
id: string;
134+
total?: number;
135+
current?: number;
136+
filename?: string;
137+
type: RsProgressType;
138+
}
139+
125140
interface UploadProgressMessage {
126141
library: string;
127-
mediaId: string;
128-
progress: number;
129-
// ... additional fields
142+
progress: RsProgress;
143+
remainingSecondes?: number;
144+
}
145+
146+
interface ConvertProgress {
147+
id: string;
148+
filename: string;
149+
convertedId?: string | null;
150+
done: boolean;
151+
percent: number;
152+
status: string;
153+
remainingSecondes?: number | null;
154+
request?: VideoConvertRequest | null;
130155
}
131156

132157
interface ConvertMessage {
133158
library: string;
134-
mediaId: string;
135-
progress: number;
136-
status: string;
159+
progress: ConvertProgress;
137160
}
138161

139162
// Content events
@@ -314,8 +337,9 @@ eventSource.addEventListener('library-status', (event) => {
314337
eventSource.addEventListener('convert_progress', (event) => {
315338
const data: SseEvent = JSON.parse(event.data);
316339
if ('ConvertProgress' in data) {
317-
const { mediaId, progress, status } = data.ConvertProgress;
318-
console.log(`Converting ${mediaId}: ${progress}% - ${status}`);
340+
const { library, progress } = data.ConvertProgress;
341+
const eta = progress.remainingSecondes ? `, ~${progress.remainingSecondes}s remaining` : '';
342+
console.log(`Converting in ${library}: ${progress.filename} ${(progress.percent * 100).toFixed(2)}% - ${progress.status}${eta}`);
319343
}
320344
});
321345

src/domain/media.rs

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,7 +172,7 @@ pub struct ConvertProgress {
172172
pub done: bool,
173173
pub percent: f64,
174174
pub status: RsVideoTranscodeStatus,
175-
pub estimated_remaining_seconds: Option<u64>,
175+
pub remaining_secondes: Option<u64>,
176176
pub request: Option<VideoConvertRequest>,
177177
}
178178

@@ -196,3 +196,26 @@ pub struct VideoMergeRequest {
196196
pub id: String,
197197
pub items: Vec<VideoMergeItem>,
198198
}
199+
200+
#[cfg(test)]
201+
mod tests {
202+
use super::*;
203+
204+
#[test]
205+
fn convert_progress_serializes_upload_eta_field_name() {
206+
let progress = ConvertProgress {
207+
id: "convert-1".to_string(),
208+
filename: "video.mp4".to_string(),
209+
converted_id: None,
210+
done: false,
211+
percent: 0.5,
212+
status: RsVideoTranscodeStatus::Processing,
213+
remaining_secondes: Some(42),
214+
request: None,
215+
};
216+
217+
let json = serde_json::to_string(&progress).unwrap();
218+
assert!(json.contains("remainingSecondes"));
219+
assert!(!json.contains("estimatedRemainingSeconds"));
220+
}
221+
}

src/model/medias.rs

Lines changed: 28 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2682,6 +2682,8 @@ impl ModelController {
26822682
// Spawn background task
26832683
tokio::spawn(async move {
26842684
let mut poll_interval = Duration::from_secs(2);
2685+
let start_time = std::time::Instant::now();
2686+
let mut remaining = None;
26852687

26862688
loop {
26872689
sleep(poll_interval).await;
@@ -2701,6 +2703,14 @@ impl ModelController {
27012703

27022704
let done = is_terminal;
27032705
let percent = current_status.progress;
2706+
if percent > 1.0 {
2707+
let percent_ratio = percent / 100.0;
2708+
if percent_ratio > 0.01 {
2709+
let duration_from_start = start_time.elapsed().as_secs_f32();
2710+
let remaining_time = duration_from_start / percent_ratio * (1.0 - percent_ratio);
2711+
remaining = Some(remaining_time as u64);
2712+
}
2713+
}
27042714

27052715
let message = ConvertMessage {
27062716
library: lib_progress.clone(),
@@ -2712,7 +2722,7 @@ impl ModelController {
27122722
status: current_status.status.clone(),
27132723
id: request_clone.id.clone(),
27142724
request: Some(request_clone.clone()),
2715-
estimated_remaining_seconds: None,
2725+
remaining_secondes: remaining,
27162726
},
27172727
};
27182728
mc_progress.send_convert_progress(message);
@@ -2745,12 +2755,17 @@ impl ModelController {
27452755
{
27462756
Ok(reader) => {
27472757
attempts += 1;
2758+
let mut converted_infos = media_progress.clone();
2759+
converted_infos.size = reader.size;
2760+
if reader.mime.is_some() {
2761+
converted_infos.mimetype = reader.mime.clone();
2762+
}
27482763

27492764
match mc_progress
27502765
.add_library_file(
27512766
&lib_progress,
27522767
&name_progress,
2753-
Some(media_progress.clone()),
2768+
Some(converted_infos),
27542769
reader.stream,
27552770
&ConnectedUser::ServerAdmin,
27562771
)
@@ -2773,6 +2788,7 @@ impl ModelController {
27732788
}
27742789
}
27752790
Err(e) => {
2791+
attempts += 1;
27762792
log_error(
27772793
crate::tools::log::LogServiceType::Source,
27782794
format!("Failed to create reader: {:?}", e),
@@ -2781,6 +2797,9 @@ impl ModelController {
27812797
"Failed to create reader: {:?}",
27822798
e
27832799
);
2800+
if attempts >= max_attempts {
2801+
break Err(e);
2802+
}
27842803
}
27852804
}
27862805
};
@@ -2981,7 +3000,7 @@ impl ModelController {
29813000
status: RsVideoTranscodeStatus::Processing,
29823001
id: request_progress.id.clone(),
29833002
request: Some(request_progress.clone()),
2984-
estimated_remaining_seconds: remaining,
3003+
remaining_secondes: remaining,
29853004
},
29863005
};
29873006
mc_progress.send_convert_progress(message);
@@ -3145,7 +3164,7 @@ impl ModelController {
31453164
done: false,
31463165
percent: 0.0,
31473166
status: RsVideoTranscodeStatus::Pending,
3148-
estimated_remaining_seconds: None,
3167+
remaining_secondes: None,
31493168
request: None,
31503169
},
31513170
};
@@ -3178,7 +3197,7 @@ impl ModelController {
31783197
status: RsVideoTranscodeStatus::Processing,
31793198
id: request_id.clone(),
31803199
request: None,
3181-
estimated_remaining_seconds: None,
3200+
remaining_secondes: None,
31823201
},
31833202
};
31843203
mc_progress.send_convert_progress(message);
@@ -3239,7 +3258,7 @@ impl ModelController {
32393258
status: RsVideoTranscodeStatus::Processing,
32403259
id: request_id.clone(),
32413260
request: None,
3242-
estimated_remaining_seconds: None,
3261+
remaining_secondes: None,
32433262
},
32443263
};
32453264
mc_progress.send_convert_progress(message);
@@ -3283,7 +3302,7 @@ impl ModelController {
32833302
done: false,
32843303
percent: 0.95,
32853304
status: RsVideoTranscodeStatus::Processing,
3286-
estimated_remaining_seconds: None,
3305+
remaining_secondes: None,
32873306
request: None,
32883307
},
32893308
});
@@ -3356,7 +3375,7 @@ impl ModelController {
33563375
done: true,
33573376
percent: 1.0,
33583377
status: RsVideoTranscodeStatus::Completed,
3359-
estimated_remaining_seconds: None,
3378+
remaining_secondes: None,
33603379
request: None,
33613380
},
33623381
});
@@ -3373,7 +3392,7 @@ impl ModelController {
33733392
done: true,
33743393
percent: 0.0,
33753394
status: RsVideoTranscodeStatus::Failed,
3376-
estimated_remaining_seconds: None,
3395+
remaining_secondes: None,
33773396
request: None,
33783397
},
33793398
});

src/model/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,7 +109,7 @@ impl VideoConvertQueueElement {
109109
done: false,
110110
percent: 0f64,
111111
status: RsVideoTranscodeStatus::Pending,
112-
estimated_remaining_seconds: None,
112+
remaining_secondes: None,
113113
request: Some(request.clone()),
114114
},
115115
request,

src/model/plugins/video_convert_plugin.rs

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -81,9 +81,28 @@ impl ModelController {
8181

8282
let plugin_with_credentials= self.get_plugin_with_credential(plugin_id).await?;
8383
let temp_url =ModelController::get_temporary_read_url(library_id, media_id, Some(21600)).await?;
84+
let media = self
85+
.store
86+
.get_library_store(library_id)?
87+
.get_media(media_id, None)
88+
.await?;
89+
let source = if let Some(media) = media {
90+
RsRequest {
91+
url: temp_url,
92+
size: media.item.size,
93+
mime: Some(media.item.mimetype),
94+
filename: Some(media.item.name),
95+
..Default::default()
96+
}
97+
} else {
98+
RsRequest {
99+
url: temp_url,
100+
..Default::default()
101+
}
102+
};
84103
let jobrequest = RsVideoTranscodeJobPluginRequest {
85104
job: RsVideoTranscodeJob {
86-
source: RsRequest { url: temp_url, ..Default::default() },
105+
source,
87106
request
88107
},
89108
credentials: plugin_with_credentials.credential.clone().map(|r| r.into()).unwrap_or_default(),
@@ -224,4 +243,4 @@ impl PluginManager {
224243
return Err(crate::error::RsError::PluginNotFound(plugin_with_cred.plugin.id));
225244
}
226245
}
227-
}
246+
}

src/plugins/sources/async_reader_progress.rs

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,13 @@ use std::task::{Context, Poll};
55
use bytes::Bytes;
66

77
use crate::domain::progress::RsProgress;
8+
9+
const PROGRESS_UPDATE_BYTES: usize = 4 * 1024 * 1024;
10+
811
pub struct ProgressReader<R> {
912
pub inner: R,
1013
pub bytes_read: usize,
14+
pub bytes_reported: usize,
1115
pub progress_template: RsProgress,
1216
pub sender: Sender<RsProgress>
1317
}
@@ -21,7 +25,7 @@ impl<R> Drop for ProgressReader<R> {
2125
new_progress.total = Some(self.bytes_read as u64);
2226

2327
tokio::spawn(async move {
24-
sender.send(new_progress).await.unwrap();
28+
let _ = sender.send(new_progress).await;
2529
});
2630
}
2731
}
@@ -32,7 +36,8 @@ impl<R> ProgressReader<R> {
3236
inner,
3337
progress_template,
3438
sender,
35-
bytes_read: 0
39+
bytes_read: 0,
40+
bytes_reported: 0
3641
}
3742
}
3843
}
@@ -43,20 +48,19 @@ impl<R: AsyncRead + Unpin> AsyncRead for ProgressReader<R> {
4348
cx: &mut Context<'_>,
4449
buf: &mut ReadBuf<'_>,
4550
) -> Poll<io::Result<()>> {
51+
let before = buf.filled().len();
4652
let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
4753
if let Poll::Ready(Ok(())) = poll {
48-
//println!("progress: {}", self.bytes_read);
49-
self.bytes_read += buf.filled().len();
54+
let read = buf.filled().len().saturating_sub(before);
55+
self.bytes_read += read;
5056
let mut new_progress = self.progress_template.clone();
5157
new_progress.current = Some(self.bytes_read as u64);
52-
if buf.filled().len() > 0 {
53-
let sender = self.sender.clone();
54-
tokio::spawn(async move {
55-
sender.send(new_progress).await.unwrap();
56-
});
58+
if read > 0 && self.bytes_read.saturating_sub(self.bytes_reported) >= PROGRESS_UPDATE_BYTES {
59+
self.bytes_reported = self.bytes_read;
60+
let _ = self.sender.try_send(new_progress);
5761
}
5862

5963
}
6064
poll
6165
}
62-
}
66+
}

src/plugins/sources/mod.rs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,8 @@ pub mod async_reader_progress;
2626

2727
pub type AsyncReadPinBox = Pin<Box<dyn AsyncRead + Send + Sync>>;
2828

29+
const STREAM_BUFFER_SIZE: usize = 4 * 1024 * 1024;
30+
2931
pub trait AsyncSeekableWrite: AsyncWrite + AsyncSeek + Send {}
3032

3133
impl<T> AsyncSeekableWrite for T where T: AsyncWrite + AsyncSeek + Send {}
@@ -327,7 +329,7 @@ impl SourceRead {
327329
match self {
328330
SourceRead::Stream(reader) => {
329331
let headers = reader.hearders().map_err(|_| Error::UnableToFormatHeaders)?;
330-
let stream = ReaderStream::new(reader.stream);
332+
let stream = ReaderStream::with_capacity(reader.stream, STREAM_BUFFER_SIZE);
331333
let body = Body::from_stream(stream);
332334
let status = if reader.range.is_some() { axum::http::StatusCode::PARTIAL_CONTENT } else { axum::http::StatusCode::OK };
333335
Ok((status, headers, body).into_response())
@@ -487,4 +489,4 @@ mod tests {
487489
assert_eq!(parsed.end, Some(1023), "test end parsing");
488490
Ok(())
489491
}
490-
}
492+
}

0 commit comments

Comments
 (0)