fix(security): harden sidecar and media processing lifecycles

This commit is contained in:
NimBold
2026-06-17 09:45:31 +03:30
parent 7c7317adc9
commit 58f4a8a14d
8 changed files with 135 additions and 21 deletions
+4
View File
@@ -7,6 +7,9 @@ use ts_rs::TS;
#[ts(export, export_to = "../../src/bindings/")]
pub enum DownloadStatus {
Downloading,
/// Post-download media processing such as yt-dlp/ffmpeg merging or
/// extraction. The queue permit is still held.
Processing,
Paused,
Completed,
Failed,
@@ -20,6 +23,7 @@ impl DownloadStatus {
pub fn as_str(self) -> &'static str {
match self {
Self::Downloading => "downloading",
Self::Processing => "processing",
Self::Paused => "paused",
Self::Completed => "completed",
Self::Failed => "failed",
+87 -6
View File
@@ -39,6 +39,61 @@ pub struct MediaMetadata {
pub formats: Vec<MediaFormat>,
}
fn is_media_processing_line(line: &str) -> bool {
let lower = line.to_lowercase();
lower.contains("[merger]")
|| lower.contains("[extractaudio]")
|| lower.contains("[ffmpeg]")
|| lower.contains("[videoconvertor]")
|| lower.contains("[fixup")
|| lower.contains("merging formats")
|| lower.contains("post-process")
}
async fn cleanup_media_processing_artifacts(out_path: &std::path::Path) {
let Some(parent) = out_path.parent() else {
return;
};
let Some(base_name) = out_path.file_name().and_then(|name| name.to_str()) else {
return;
};
let base_stem = out_path
.file_stem()
.and_then(|name| name.to_str())
.unwrap_or(base_name);
let _ = tokio::fs::remove_file(out_path).await;
let Ok(mut entries) = tokio::fs::read_dir(parent).await else {
return;
};
while let Ok(Some(entry)) = entries.next_entry().await {
let path = entry.path();
if path == out_path {
continue;
}
let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
if !name.starts_with(base_name) && !name.starts_with(base_stem) {
continue;
}
let yt_dlp_format_fragment = name
.strip_prefix(base_stem)
.and_then(|suffix| suffix.strip_prefix(".f"))
.and_then(|suffix| suffix.chars().next())
.is_some_and(|ch| ch.is_ascii_digit());
let looks_like_media_temp = name.contains(".part")
|| name.contains(".ytdl")
|| name.contains(".temp")
|| name.contains(".tmp")
|| yt_dlp_format_fragment;
if looks_like_media_temp {
let _ = tokio::fs::remove_file(path).await;
}
}
}
@@ -693,6 +748,7 @@ pub(crate) async fn start_media_download_internal(
let mut strike = 0_usize;
let mut terminal_failure = false;
let mut processing_started = false;
'retry: while strike <= MAX_RETRIES {
let mut cmd = app_handle.shell().sidecar("yt-dlp").map_err(|e| e.to_string())?
@@ -772,7 +828,10 @@ pub(crate) async fn start_media_download_internal(
tokio::select! {
_ = cancel_rx.changed() => {
let _ = child.kill();
return Ok(());
if processing_started {
cleanup_media_processing_artifacts(&out_path).await;
}
return Err(crate::queue::MEDIA_RUN_CANCELLED.to_string());
}
event = rx.recv() => {
match event {
@@ -819,6 +878,23 @@ pub(crate) async fn start_media_download_internal(
}
Some(tauri_plugin_shell::process::CommandEvent::Stderr(line_bytes)) => {
let line = String::from_utf8_lossy(&line_bytes);
if !processing_started && is_media_processing_line(&line) {
processing_started = true;
let _ = app_handle.emit(
"download-state",
DownloadStateEvent::new(
id,
crate::ipc::DownloadStatus::Processing,
),
);
let _ = app_handle.emit("download-progress", DownloadProgressEvent {
id: id.to_string(),
fraction: 1.0,
speed: "Processing".to_string(),
eta: "-".to_string(),
size: None,
});
}
let lower = line.to_lowercase();
if lower.contains("error") || lower.contains("critical") {
log::error!("yt-dlp stderr [{}]: {}", id, line.trim());
@@ -911,8 +987,7 @@ async fn pause_download(
) -> Result<(), String> {
println!("pause_download called for id: {}", id);
// Release the concurrency slot FIRST, before signaling the sidecar to die.
state.queue_manager.release_permit(&id).await;
let active_kind = state.queue_manager.active_kind(&id).await;
state.queue_manager.remove_from_pending(&id).await;
// Emit the paused state.
@@ -938,7 +1013,11 @@ async fn pause_download(
.send(download::DownloadCmd::Pause(download_id))
.await;
}
state.download_coordinator.pause_media(id).await
let media_result = state.download_coordinator.pause_media(id.clone()).await;
if !matches!(active_kind, Some(crate::queue::TaskKind::Media)) {
state.queue_manager.release_permit(&id).await;
}
media_result
}
#[tauri::command]
@@ -950,9 +1029,8 @@ async fn remove_download(
) -> Result<(), String> {
println!("remove_download called for id: {}", id);
// Remove from the queue (pending or active) and free the slot.
let active_kind = state.queue_manager.active_kind(&id).await;
state.queue_manager.remove_from_pending(&id).await;
state.queue_manager.release_permit(&id).await;
use tauri::Emitter;
let _ = app_handle.emit(
@@ -967,6 +1045,9 @@ async fn remove_download(
.await?;
}
state.download_coordinator.pause_media(id.clone()).await?;
if !matches!(active_kind, Some(crate::queue::TaskKind::Media)) {
state.queue_manager.release_permit(&id).await;
}
if let Some(path) = filepath {
if !path.is_empty() {
+14 -2
View File
@@ -12,6 +12,7 @@ use log;
/// Default capacity when no setting is read yet.
pub const DEFAULT_MAX_CONCURRENT: usize = 3;
pub const MEDIA_RUN_CANCELLED: &str = "__firelink_media_run_cancelled__";
/// Outcome of an aria2 completion that arrived before its gid was stored.
/// Carries the outcome so the correct state emit survives the race.
@@ -84,6 +85,7 @@ pub struct QueueManager<R: tauri::Runtime = tauri::Wry> {
pending: Mutex<VecDeque<QueuedTask>>,
semaphore: Arc<Semaphore>,
active_permits: Mutex<HashMap<String, OwnedSemaphorePermit>>,
active_kinds: Mutex<HashMap<String, TaskKind>>,
target_capacity: AtomicUsize,
slots_to_retire: AtomicUsize,
notify: Notify,
@@ -126,6 +128,7 @@ impl<R: tauri::Runtime> QueueManager<R> {
pending: Mutex::new(VecDeque::new()),
semaphore: Arc::new(Semaphore::new(capacity)),
active_permits: Mutex::new(HashMap::new()),
active_kinds: Mutex::new(HashMap::new()),
target_capacity: AtomicUsize::new(capacity),
slots_to_retire: AtomicUsize::new(0),
notify: Notify::new(),
@@ -178,10 +181,15 @@ impl<R: tauri::Runtime> QueueManager<R> {
.insert(id.to_string(), permit);
}
pub async fn active_kind(&self, id: &str) -> Option<TaskKind> {
self.active_kinds.lock().await.get(id).cloned()
}
/// Release the permit parked under `id`, if any. Idempotent. Wakes the
/// dispatcher so a freed slot is claimed promptly.
pub async fn release_permit(&self, id: &str) {
let removed = self.active_permits.lock().await.remove(id).is_some();
self.active_kinds.lock().await.remove(id);
if removed {
self.notify.notify_one();
}
@@ -298,6 +306,7 @@ impl<R: tauri::Runtime> QueueManager<R> {
// aria2's RPC returns instantly, so the permit must outlive the
// dispatch_one call. Media/Native runners release on exit.
self.park_permit(&id, permit).await;
self.active_kinds.lock().await.insert(id.clone(), task.kind.clone());
self.emit_state(&id, DownloadStatus::Downloading);
match task.kind {
@@ -350,6 +359,7 @@ impl<R: tauri::Runtime> QueueManager<R> {
Ok(()) => {
self.emit_state(id, DownloadStatus::Completed);
}
Err(error) if error == MEDIA_RUN_CANCELLED => {}
Err(error) => {
self.emit_failed(id, error);
}
@@ -703,7 +713,7 @@ impl SidecarSpawner for ProductionSpawner {
.register_media(id.to_string())
.await
.map_err(|e| e)?;
crate::start_media_download_internal(
let outcome = crate::start_media_download_internal(
self.app_handle.clone(),
id,
payload.url.clone(),
@@ -720,7 +730,9 @@ impl SidecarSpawner for ProductionSpawner {
payload.max_tries,
&mut cancel_rx,
)
.await
.await;
let _ = state.download_coordinator.finish_media(id.to_string()).await;
outcome
}
async fn run_native(&self, id: &str, payload: &SpawnPayload) -> Result<(), String> {